Skip to main content

free_agent/
message.rs

1//! What actors say to one another: the [`Message`] trait and the request
2//! and reply envelopes that carry a message and its answer.
3
4use anyhow::{Context, Result};
5use serde::{Serialize, de::DeserializeOwned};
6use std::fmt::Debug;
7use std::sync::atomic::{AtomicU64, Ordering};
8use tokio::sync::oneshot;
9
10/// What one episode's actors say to one another.
11///
12/// An episode carries one message type, and every actor in it sends and
13/// receives that type. A message is cloned for every recipient of a
14/// statement or a request, travels between actors on a multi-threaded
15/// runtime, and is written to the log and read back from it, which is what
16/// the bounds say. Implementing it takes one empty
17/// line:
18///
19/// ```
20/// # use free_agent::Message;
21/// # use serde::{Deserialize, Serialize};
22/// #[derive(Debug, Clone, Serialize, Deserialize)]
23/// struct Note(String);
24/// impl Message for Note {}
25/// ```
26pub trait Message: Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static {}
27
28/// Tells a reply from the others arriving on the same channel.
29type RequestId = u64;
30
31/// Request ids are unique across the process.
32static NEXT_REQUEST: AtomicU64 = AtomicU64::new(0);
33
34/// A message whose sender is waiting for an answer, and the channel the
35/// answer goes back on.
36#[derive(Debug)]
37pub(crate) struct Request<M: Message> {
38    id: RequestId,
39    message: M,
40    reply_to: oneshot::Sender<Reply<M>>,
41}
42
43impl<M: Message> Request<M> {
44    /// A request carrying `message`, and the channel its reply arrives on.
45    pub(crate) fn new(message: M) -> (Self, oneshot::Receiver<Reply<M>>) {
46        let (reply_to, reply) = oneshot::channel();
47        let request = Self {
48            id: NEXT_REQUEST.fetch_add(1, Ordering::Relaxed),
49            message,
50            reply_to,
51        };
52        (request, reply)
53    }
54
55    /// Answer this request with `messages`, consuming it: a request is
56    /// answered once.
57    ///
58    /// # Errors
59    ///
60    /// Fails if the asker has stopped waiting.
61    pub(crate) fn reply(self, messages: Vec<M>) -> Result<()> {
62        self.reply_to
63            .send(Reply {
64                id: self.id,
65                messages,
66            })
67            .ok()
68            .context("the asker stopped waiting")
69    }
70
71    /// What is being asked.
72    pub(crate) fn message(&self) -> &M {
73        &self.message
74    }
75}
76
77/// The answer to a [`Request`].
78#[derive(Debug)]
79pub(crate) struct Reply<M: Message> {
80    /// The request this answers.
81    #[allow(dead_code)]
82    id: RequestId,
83    /// What the recipient answered. An empty answer is an acknowledgment.
84    pub(crate) messages: Vec<M>,
85}