compio_actor\process_group/
mod.rs1mod 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
18pub struct ProcessGroup<M: Message> {
20 inner: Arc<GroupInner<M>>,
21}
22
23impl<M: Message> ProcessGroup<M> {
24 pub fn new() -> Self {
26 Self::with_strategy(Strategy::default())
27 }
28
29 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 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 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 pub fn len(&self) -> usize {
98 self.inner.state.lock().unwrap().members.len()
99 }
100
101 pub fn is_empty(&self) -> bool {
103 self.len() == 0
104 }
105}
106
107impl<M: Message, R: Message> ProcessGroup<Call<M, R>> {
108 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#[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 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}