behavior_actors/operations/
configuration.rs1#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
12pub struct ConfigurationVersion(pub u64);
13
14#[derive(Debug, Clone, PartialEq, Eq)]
16pub enum ConfigurationState<C> {
17 Unconfigured,
19 Configured {
21 version: ConfigurationVersion,
23 value: C,
25 },
26}
27
28pub enum ConfigurationMessage<C, Route> {
30 Apply {
32 version: ConfigurationVersion,
34 value: C,
36 },
37 Query {
39 reply_to: Route,
41 },
42}
43
44#[derive(Debug, Error, Clone, PartialEq, Eq)]
46pub enum ConfigurationError<C> {
47 #[error("configuration candidate is stale")]
49 Stale {
50 proposed: ConfigurationVersion,
52 current: ConfigurationVersion,
54 value: C,
56 },
57 #[error("configuration candidate conflicts at the committed version")]
59 ConflictingVersion {
60 version: ConfigurationVersion,
62 value: C,
64 },
65}
66
67pub struct Configuration<
79 A: Address,
80 C,
81 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
82> {
83 state: ConfigurationState<C>,
84 marker: core::marker::PhantomData<fn() -> (A, Route)>,
85}
86
87impl<A, C, Route> Configuration<A, C, Route>
88where
89 A: Address,
90 C: Clone + Eq,
91 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
92{
93 #[must_use]
95 pub const fn new() -> Self {
96 Self {
97 state: ConfigurationState::Unconfigured,
98 marker: core::marker::PhantomData,
99 }
100 }
101
102 #[must_use]
104 pub const fn state(&self) -> &ConfigurationState<C> {
105 &self.state
106 }
107
108 fn apply(
109 &mut self,
110 version: ConfigurationVersion,
111 value: C,
112 ) -> Result<(), ConfigurationError<C>> {
113 let ConfigurationState::Configured {
114 version: current,
115 value: committed,
116 } = &self.state
117 else {
118 self.state = ConfigurationState::Configured { version, value };
119 return Ok(());
120 };
121 if version < *current {
122 return Err(ConfigurationError::Stale {
123 proposed: version,
124 current: *current,
125 value,
126 });
127 }
128 if version == *current {
129 if value == *committed {
130 return Ok(());
131 }
132 return Err(ConfigurationError::ConflictingVersion { version, value });
133 }
134 self.state = ConfigurationState::Configured { version, value };
135 Ok(())
136 }
137}
138
139impl<A, C, Route> Default for Configuration<A, C, Route>
140where
141 A: Address,
142 C: Clone + Eq,
143 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
144{
145 fn default() -> Self {
146 Self::new()
147 }
148}
149
150impl<A, C, Route> BehaviorBase for Configuration<A, C, Route>
151where
152 A: Address,
153 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
154{
155 type Base = Self;
156 fn base(&self) -> &Self {
157 self
158 }
159}
160
161impl<A, C, Route> behavior::Protocol for Configuration<A, C, Route>
162where
163 A: Address,
164 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
165{
166 type Addr = A;
167 type Msg = ConfigurationMessage<C, Route>;
168}
169
170impl<A, C, Route> Behavior for Configuration<A, C, Route>
171where
172 A: Address,
173 C: Clone + Eq,
174 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ConfigurationState<C>>>,
175 Route::Sends: behavior::SendsFor<User<A, ConfigurationMessage<C, Route>>>,
176{
177 type Protocol = Self;
178 type Event = User<A, behavior::BehaviorMessage<Self>>;
179 type Sends = Route::Sends;
180 type Ph = Never;
181 type Error = ConfigurationError<C>;
182 type Birth = NoBirths;
183
184 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
185 match event.message {
186 ConfigurationMessage::Apply { version, value } => {
187 self.apply(version, value)?;
188 Ok(Actions::cont())
189 }
190 ConfigurationMessage::Query { reply_to } => {
191 Ok(Actions::send(reply_to.deliver(self.state.clone())))
192 }
193 }
194 }
195}
196
197#[cfg(test)]
198mod tests {
199 use crate::Activate as _;
200 use behavior::MailAddr;
201
202 use super::*;
203
204 struct Reply;
205 impl behavior::Protocol for Reply {
206 type Addr = MailAddr;
207 type Msg = ConfigurationState<u8>;
208 }
209
210 impl Behavior for Reply {
211 type Protocol = Self;
212 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
213 type Sends = Vec<Never>;
214 type Ph = Never;
215 type Error = Never;
216 type Birth = NoBirths;
217 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
218 Ok(Actions::cont())
219 }
220 }
221
222 type Subject = Configuration<MailAddr, u8, Recipient<Reply>>;
223
224 #[test]
225 fn stale_and_conflicting_candidates_return_ownership_atomically() {
226 let mut subject = (Subject::new()).initialize().unwrap().behavior;
227 let applied = subject
228 .receive(
229 MailAddr(9),
230 ConfigurationMessage::Apply {
231 version: ConfigurationVersion(2),
232 value: 20,
233 },
234 )
235 .unwrap();
236 assert!(applied.sends.is_empty());
237 assert!(applied.creates.is_empty());
238 assert_eq!(applied.become_, behavior::Step::Continue);
239 let rejection = subject.receive(
240 MailAddr(9),
241 ConfigurationMessage::Apply {
242 version: ConfigurationVersion(1),
243 value: 10,
244 },
245 );
246 assert!(matches!(
247 rejection,
248 Err(ConfigurationError::Stale { value: 10, .. })
249 ));
250 let rejection = subject.receive(
251 MailAddr(9),
252 ConfigurationMessage::Apply {
253 version: ConfigurationVersion(2),
254 value: 21,
255 },
256 );
257 assert!(matches!(
258 rejection,
259 Err(ConfigurationError::ConflictingVersion { value: 21, .. })
260 ));
261 assert_eq!(
262 subject.state(),
263 &ConfigurationState::Configured {
264 version: ConfigurationVersion(2),
265 value: 20
266 }
267 );
268 }
269
270 #[test]
271 fn query_reports_unconfigured_and_configured_as_distinct_states() {
272 let mut subject = (Subject::new()).initialize().unwrap().behavior;
273 let initial = subject
274 .receive(
275 MailAddr(9),
276 ConfigurationMessage::Query {
277 reply_to: Recipient::global(MailAddr(1)),
278 },
279 )
280 .unwrap();
281 assert!(matches!(
282 initial.sends[0].message,
283 ConfigurationState::Unconfigured
284 ));
285 }
286}