Skip to main content

compio_actor/process_group/
mod.rs

1//! Typed actor groups with configurable routing.
2
3mod strategy;
4
5use std::{
6    fmt,
7    sync::{Arc, Mutex, Weak},
8};
9
10#[doc(inline)]
11pub use strategy::Strategy;
12
13use crate::{
14    Broker, Call, Message,
15    mailbox::{CallError, DeliverError, call_with},
16};
17
18/// A group of actors that share messages using a routing [`Strategy`].
19pub struct ProcessGroup<M: Message> {
20    inner: Arc<GroupInner<M>>,
21}
22
23impl<M: Message> ProcessGroup<M> {
24    /// Creates an empty process group.
25    pub fn new() -> Self {
26        Self::with_strategy(Strategy::default())
27    }
28
29    /// Creates an empty process group with a routing strategy.
30    pub fn with_strategy(strategy: Strategy) -> Self {
31        Self {
32            inner: Arc::new(GroupInner {
33                state: Mutex::new(GroupState {
34                    next_id: 0,
35                    cursor: 0,
36                    members: Vec::new(),
37                    strategy,
38                }),
39            }),
40        }
41    }
42
43    /// Adds a broker until the returned membership is dropped.
44    pub fn join(&self, broker: Broker<M>) -> Membership<M> {
45        let mut state = self.inner.state.lock().unwrap();
46        let id = state.next_id;
47        state.next_id = state.next_id.wrapping_add(1);
48        state.members.push(Member { id, broker });
49        Membership {
50            id,
51            group: Arc::downgrade(&self.inner),
52        }
53    }
54
55    /// Routes a message to the next available member.
56    pub fn send(&self, mut message: M) -> Result<(), DeliverError<M>> {
57        let mut state = self.inner.state.lock().unwrap();
58        let attempts = state.members.len();
59        let mut attempted = 0;
60        let mut saw_full = false;
61        let mut index = match std::num::NonZeroUsize::new(state.members.len()) {
62            Some(members) => {
63                let strategy = state.strategy;
64                strategy.select(&mut state.cursor, members)
65            }
66            None => return Err(DeliverError::Closed(message)),
67        };
68
69        while attempted < attempts && !state.members.is_empty() {
70            attempted += 1;
71
72            match state.members[index].broker.send(message) {
73                Ok(()) => return Ok(()),
74                Err(DeliverError::Full(returned)) => {
75                    saw_full = true;
76                    message = returned;
77                    index = (index + 1) % state.members.len();
78                }
79                Err(DeliverError::Closed(returned)) => {
80                    message = returned;
81                    state.members.remove(index);
82                    if !state.members.is_empty() {
83                        index %= state.members.len();
84                    }
85                }
86            }
87        }
88
89        if saw_full {
90            Err(DeliverError::Full(message))
91        } else {
92            Err(DeliverError::Closed(message))
93        }
94    }
95
96    /// Returns the number of registered members.
97    pub fn len(&self) -> usize {
98        self.inner.state.lock().unwrap().members.len()
99    }
100
101    /// Returns whether the group has no members.
102    pub fn is_empty(&self) -> bool {
103        self.len() == 0
104    }
105}
106
107impl<M: Message, R: Message> ProcessGroup<Call<M, R>> {
108    /// Routes a request and waits for the selected actor's reply.
109    pub async fn call(&self, message: M) -> Result<R, CallError<M>> {
110        call_with(message, |call| self.send(call)).await
111    }
112}
113
114impl<M: Message> Clone for ProcessGroup<M> {
115    fn clone(&self) -> Self {
116        Self {
117            inner: self.inner.clone(),
118        }
119    }
120}
121
122impl<M: Message> Default for ProcessGroup<M> {
123    fn default() -> Self {
124        Self::new()
125    }
126}
127
128impl<M: Message> fmt::Debug for ProcessGroup<M> {
129    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
130        f.debug_struct("ProcessGroup")
131            .field("members", &self.len())
132            .finish()
133    }
134}
135
136/// An actor's membership in a [`ProcessGroup`].
137#[must_use = "dropping the membership removes the actor from the process group"]
138pub struct Membership<M: Message> {
139    id: u64,
140    group: Weak<GroupInner<M>>,
141}
142
143impl<M: Message> Membership<M> {
144    /// Removes the actor from the group.
145    pub fn leave(self) {}
146}
147
148impl<M: Message> Drop for Membership<M> {
149    fn drop(&mut self) {
150        let Some(group) = self.group.upgrade() else {
151            return;
152        };
153        let mut state = group.state.lock().unwrap();
154        if let Some(index) = state.members.iter().position(|member| member.id == self.id) {
155            state.members.remove(index);
156        }
157    }
158}
159
160impl<M: Message> fmt::Debug for Membership<M> {
161    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
162        f.debug_struct("Membership").field("id", &self.id).finish()
163    }
164}
165
166struct GroupInner<M: Message> {
167    state: Mutex<GroupState<M>>,
168}
169
170struct GroupState<M: Message> {
171    next_id: u64,
172    cursor: usize,
173    members: Vec<Member<M>>,
174    strategy: Strategy,
175}
176
177struct Member<M: Message> {
178    id: u64,
179    broker: Broker<M>,
180}