Skip to main content

free_agent/
actor.rs

1//! Actors in the sense of the [actor model]: each one keeps its own state,
2//! takes messages from a mailbox one at a time, and reaches other actors
3//! only by sending them messages.
4//!
5//! Here an actor is a [`Behavior`] driven by a mailbox. The behavior
6//! reaches the rest of the episode through its [`Context`], which an
7//! [`Episode`](crate::Episode) builds from the actor's [`ActorInit`].
8//!
9//! [actor model]: https://en.wikipedia.org/wiki/Actor_model
10
11use 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
23/// An actor's name, unique within an episode.
24pub type ActorId = String;
25
26/// Builds a [`Behavior`] from the [`Context`] it will own.
27pub type Builder<B> =
28    Box<dyn FnOnce(Context<<B as Behavior>::Message, <B as Behavior>::Log>) -> B + Send>;
29
30/// Everything an [`Episode`](crate::Episode) needs to build one actor: how
31/// it behaves and whom it may reach.
32pub struct ActorInit<B: Behavior> {
33    /// Builds the behavior once its [`Context`] exists, which is when the
34    /// episode has opened every actor's channels.
35    pub behavior: Builder<B>,
36    /// The actors this one may send to and request from.
37    pub can_send_to: HashSet<ActorId>,
38    /// The actors this one may stop.
39    pub can_shut_down: HashSet<ActorId>,
40    /// Whether this actor gets a copy of the episode's logger.
41    pub has_logger: bool,
42}
43
44/// A running actor: its behavior, its mailbox, and the two one-time signals
45/// it exchanges with the episode on the way up.
46pub(crate) struct Actor<B: Behavior> {
47    /// What this Actor does at each point in its life. It owns the
48    /// [`Context`] through which it reaches other actors.
49    pub(crate) behavior: B,
50    /// This Actor's one-time signal to the episode that it is initialized
51    /// and running its message loop. Taken when it is sent.
52    pub(crate) ready: Option<oneshot::Sender<()>>,
53    /// The episode's one-time signal that every actor is running.
54    pub(crate) start: oneshot::Receiver<()>,
55    /// The channel on which this Actor receives incoming Messages.
56    pub(crate) mailbox: UnboundedReceiver<Envelope<B::Message>>,
57}
58
59impl<B: Behavior> Actor<B> {
60    /// An actor's whole life, in four phases.
61    ///
62    /// 1. Initialize, and tell the episode this actor is ready.
63    /// 2. Wait for the episode's start signal, which comes once every actor
64    ///    is ready, and make the opening move.
65    /// 3. Deliver each envelope in the mailbox to the behavior, one at a
66    ///    time, until shut down.
67    /// 4. Clean up.
68    ///
69    /// Phases 2 and 3 overlap: mail that arrives before the start signal is
70    /// delivered as it comes. An error from the behavior at any phase ends
71    /// the actor with that error.
72    ///
73    /// Shutdown cuts any phase short. A step in progress is dropped at its
74    /// next await, mail still in the mailbox stays there, and cleanup runs
75    /// only after a finished initialization.
76    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        // The episode may already be gone, in which case no one is waiting.
84        if let Some(ready) = self.ready.take() {
85            let _ = ready.send(());
86        }
87        let mut start_consumed = false;
88        loop {
89            // Wait for whichever happens first: shutdown, the start signal,
90            // or the next envelope. Each arm is `pattern = future => body`;
91            // the body runs with the future's output bound to the pattern.
92            // `biased` tries the arms in order, so shutdown wins a tie. The
93            // start arm drops out once it has fired, because a oneshot
94            // receiver panics if polled again.
95            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                    // This actor holds a sender to its own mailbox, so the
104                    // mailbox outlives the loop.
105                    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    /// Make the opening move, now that the episode has signaled that every
117    /// actor is running. An error in place of the signal means the episode
118    /// is gone, and there is nothing to open.
119    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    /// Hand `envelope` to the behavior: either a statement to receive or a
135    /// request to answer and reply to.
136    ///
137    /// The message loop watches for shutdown only between envelopes, so each
138    /// step races shutdown on its own here.
139    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                // The asker may have stopped waiting, which is no fault of
157                // this actor.
158                let _ = request.reply(answered?);
159            }
160        }
161        Ok(Continue(()))
162    }
163}
164
165/// Run `step` to completion, unless `shutdown` fires first, in which case
166/// the step is dropped where it stands.
167async 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
178/// The ways out of an actor: what a [`Behavior`] may do besides answer.
179///
180/// A behavior owns one of these, and the actor keeps the mailbox, so messages
181/// reach the behavior one step at a time.
182///
183/// `L` is what this actor logs. It defaults to the message type, which is
184/// what most actors log.
185pub struct Context<M: Message, L = M> {
186    /// This actor's name, as the other actors know it.
187    pub id: ActorId,
188    /// The episode this actor is in, stamped on everything it logs.
189    pub episode: Uuid,
190    /// The sending ends of the mailboxes this actor may put something in,
191    /// its own among them.
192    pub(crate) mailboxes: HashMap<ActorId, UnboundedSender<Envelope<M>>>,
193    /// Tokens on which this Actor is shut down, and shuts down others.
194    pub(crate) shutdown: Shutdown,
195    /// Where this Actor's events go, when it has a logger.
196    pub(crate) log: Option<Logger<L>>,
197}
198
199impl<M: Message, L> Context<M, L> {
200    /// Log `payload` as an event stamped now. The event reaches the log when
201    /// this actor holds a logger and the log is listening; otherwise it is
202    /// dropped, and the actor carries on either way.
203    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    /// Send the statement `message` to every actor in `to`. An actor that
210    /// has already stopped is skipped. An actor may send to itself: the
211    /// message waits in its own mailbox for the current step to finish.
212    ///
213    /// # Errors
214    ///
215    /// Fails before anything is sent when `to` names an actor outside those
216    /// this one may send to.
217    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            // A failed send means the recipient's mailbox is gone.
224            let _ = sender.send(Envelope::Statement(message.clone()));
225        }
226        Ok(())
227    }
228
229    /// Ask every actor in `to` the same thing and collect their replies. A
230    /// recipient that has stopped, before or after receiving the request,
231    /// is left out of the result.
232    ///
233    /// # Errors
234    ///
235    /// Fails before anything is sent when `to` names an actor outside those
236    /// this one may send to, or names this actor itself, which is busy with this
237    /// very step.
238    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            // A dropped reply channel means the recipient stopped first.
263            if let Ok(reply) = reply.await {
264                replies.insert(id, reply.messages);
265            }
266        }
267        Ok(replies)
268    }
269
270    /// The sending end of the mailbox of `id`, which is this actor's own
271    /// when `id` is its own name.
272    ///
273    /// # Errors
274    ///
275    /// Fails when `id` is not an actor this one may send to.
276    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    /// Stop another actor at once, wherever it is in its step.
286    ///
287    /// # Errors
288    ///
289    /// Fails when `who` is outside the actors this one may shut down.
290    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    /// Shut this actor down. The step that calls this is abandoned at its
304    /// next await, so a behavior with last words says them first. The actor
305    /// then stops, leaving whatever is in its mailbox there.
306    pub fn shutdown(&self) {
307        self.shutdown.mine.cancel();
308    }
309}
310
311/// What an actor does at each point in its life. It initializes, makes an
312/// opening move once every actor is running, then maps what it receives to
313/// what it does: a statement goes to [`receive`](Behavior::receive) and a
314/// request to [`answer`](Behavior::answer), whose result is the reply. On
315/// the way out it cleans up. Every method but [`context`](Behavior::context)
316/// has an empty default, so the simplest behavior is a context and nothing
317/// else. Anything else the behavior wants to say, such as who it is, goes
318/// inside its messages.
319///
320/// A behavior's state is its own: the mailbox hands it one message at a
321/// time, so each step may change that state freely. Actors run on a
322/// multi-threaded runtime, so a behavior has to be sendable between threads.
323///
324/// A behavior keeps the [`Context`] it is built with and returns it from
325/// [`context`](Behavior::context). The episode makes the context, and the
326/// [`Builder`] in the actor's [`ActorInit`] wraps the behavior around it:
327///
328/// ```
329/// # use async_trait::async_trait;
330/// # use free_agent::{ActorInit, Behavior, Context, Message};
331/// # use std::collections::HashSet;
332/// # use serde::{Deserialize, Serialize};
333/// #[derive(Debug, Clone, Serialize, Deserialize)]
334/// struct Note(String);
335/// impl Message for Note {}
336///
337/// /// Writes down everything it hears.
338/// struct Scribe {
339///     context: Context<Note>,
340/// }
341/// #[async_trait]
342/// impl Behavior for Scribe {
343///     type Message = Note;
344///     type Log = Note;
345///     fn context(&self) -> &Context<Note> {
346///         &self.context
347///     }
348///     async fn receive(&mut self, message: &Note) -> anyhow::Result<()> {
349///         self.context.log(message.clone());
350///         Ok(())
351///     }
352/// }
353///
354/// let init = ActorInit {
355///     behavior: Box::new(|context| Scribe { context }),
356///     can_send_to: HashSet::new(),
357///     can_shut_down: HashSet::new(),
358///     has_logger: true,
359/// };
360/// # let _: ActorInit<Scribe> = init;
361/// ```
362#[async_trait]
363pub trait Behavior: Send {
364    /// What this behavior sends and receives.
365    type Message: Message;
366    /// What this behavior logs. Most often the message type.
367    type Log: Serialize + DeserializeOwned + Send + 'static;
368    /// The ways out of this actor, handed to the behavior when it was built.
369    fn context(&self) -> &Context<Self::Message, Self::Log>;
370    /// Called before anything else: open a connection, say. No actor starts
371    /// until every actor has initialized.
372    async fn initialize(&mut self) -> anyhow::Result<()> {
373        Ok(())
374    }
375    /// Handle the statement `message`. By default it is ignored.
376    async fn receive(&mut self, _message: &Self::Message) -> anyhow::Result<()> {
377        Ok(())
378    }
379    /// Answer the request `message`. The result is the reply, which by
380    /// default is nothing.
381    async fn answer(&mut self, _message: &Self::Message) -> anyhow::Result<Vec<Self::Message>> {
382        Ok(vec![])
383    }
384    /// Called once every actor in the episode is running. This is where an
385    /// actor with an opening move makes it.
386    async fn start(&mut self) -> anyhow::Result<()> {
387        Ok(())
388    }
389    /// Called last, after the actor has stopped taking mail: close what
390    /// `initialize` opened.
391    async fn clean_up(&mut self) -> anyhow::Result<()> {
392        Ok(())
393    }
394}
395
396/// The cancellation tokens an actor is stopped through and stops others
397/// through.
398pub(crate) struct Shutdown {
399    /// Token other Actors cancel to shut this Actor down.
400    pub(crate) mine: CancellationToken,
401    /// Tokens this Actor cancels to shut down other Actors.
402    pub(crate) others: HashMap<ActorId, CancellationToken>,
403}
404
405/// A message and what its recipient owes for it.
406#[derive(Debug)]
407pub(crate) enum Envelope<M: Message> {
408    /// A message that does not require a reply.
409    Statement(M),
410    /// A message whose sender is waiting for reply.
411    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    /// A behavior that answers every request with the message it was sent,
433    /// after waiting for a permit from its gate if it has one, and waits at
434    /// the gate on statements too. Slow to start, it waits for a permit in
435    /// `start` as well. Broken, it fails every step.
436    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        /// Fail if broken, otherwise wait at the gate if there is one.
468        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    /// A behavior with every default: no opening move, statements ignored,
480    /// and requests answered with nothing.
481    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    /// A behavior that counts the statements it has received and answers
492    /// every request with that many notes.
493    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    /// An actor, along with what a test needs to feed it, start it, and stop
514    /// it from the outside.
515    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    /// A rig around an [`Echo`] with no gate.
529    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    /// A rig around the behavior `build` makes from its context, which may
543    /// send to itself and to `others`.
544    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        // Bob is still running and answering afterward.
719        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    /// A rig whose actor blocks in every step until the gate gives a permit.
839    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        // Let Bob take the message and block in `receive`.
855        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        // Let Bob take the start signal and block in its opening move.
870        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        // Let Bob take the start signal and settle into waiting.
884        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        // Ann gives up after a second, dropping her reply channel, and then
964        // Bob answers.
965        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}