1#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use super::ObservationVersion;
9use crate::DeliveryRoute;
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum ReadinessStatus {
14 Ready,
16 NotReady,
18}
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub enum ReadinessEvidence {
23 Unknown,
25 Observed {
27 version: ObservationVersion,
29 status: ReadinessStatus,
31 },
32}
33
34#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct DependencyReadiness<K> {
37 pub dependency: K,
39 pub evidence: ReadinessEvidence,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
45pub struct ReadinessReport<K> {
46 pub dependencies: Vec<DependencyReadiness<K>>,
48}
49
50impl<K> ReadinessReport<K> {
51 #[must_use]
53 pub fn ready(&self) -> bool {
54 self.dependencies.iter().all(|state| {
55 matches!(
56 state.evidence,
57 ReadinessEvidence::Observed {
58 status: ReadinessStatus::Ready,
59 ..
60 }
61 )
62 })
63 }
64}
65
66pub enum ReadinessMessage<K, Route> {
68 Observe {
70 dependency: K,
72 version: ObservationVersion,
74 status: ReadinessStatus,
76 },
77 Query {
79 reply_to: Route,
81 },
82}
83
84#[derive(Debug, Error, Clone, PartialEq, Eq)]
86pub enum ReadinessError<K> {
87 #[error("readiness dependency is not configured")]
89 UnknownDependency {
90 dependency: K,
92 observed: ObservationVersion,
94 status: ReadinessStatus,
96 },
97 #[error("readiness evidence is stale")]
99 Stale {
100 dependency: K,
102 observed: ObservationVersion,
104 current: ObservationVersion,
106 status: ReadinessStatus,
108 },
109 #[error("readiness evidence conflicts at the committed version")]
111 ConflictingVersion {
112 dependency: K,
114 version: ObservationVersion,
116 status: ReadinessStatus,
118 },
119}
120
121pub struct Readiness<A, K, Route>
134where
135 A: Address,
136 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ReadinessReport<K>>>,
137{
138 dependencies: Vec<DependencyReadiness<K>>,
139 marker: core::marker::PhantomData<fn() -> (A, Route)>,
140}
141
142impl<A, K, Route> Readiness<A, K, Route>
143where
144 A: Address,
145 K: Clone + Eq,
146 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ReadinessReport<K>>>,
147{
148 #[must_use]
150 pub fn new(dependencies: impl IntoIterator<Item = K>) -> Self {
151 let mut states = Vec::new();
152 for dependency in dependencies {
153 if !states
154 .iter()
155 .any(|state: &DependencyReadiness<K>| state.dependency == dependency)
156 {
157 states.push(DependencyReadiness {
158 dependency,
159 evidence: ReadinessEvidence::Unknown,
160 });
161 }
162 }
163 Self {
164 dependencies: states,
165 marker: core::marker::PhantomData,
166 }
167 }
168
169 #[must_use]
171 pub fn dependencies(&self) -> &[DependencyReadiness<K>] {
172 &self.dependencies
173 }
174
175 fn observe(
176 &mut self,
177 dependency: K,
178 version: ObservationVersion,
179 status: ReadinessStatus,
180 ) -> Result<(), ReadinessError<K>> {
181 let Some(current) = self
182 .dependencies
183 .iter_mut()
184 .find(|state| state.dependency == dependency)
185 else {
186 return Err(ReadinessError::UnknownDependency {
187 dependency,
188 observed: version,
189 status,
190 });
191 };
192 let ReadinessEvidence::Observed {
193 version: committed,
194 status: committed_status,
195 } = current.evidence
196 else {
197 current.evidence = ReadinessEvidence::Observed { version, status };
198 return Ok(());
199 };
200 if version < committed {
201 return Err(ReadinessError::Stale {
202 dependency,
203 observed: version,
204 current: committed,
205 status,
206 });
207 }
208 if version == committed {
209 if committed_status == status {
210 return Ok(());
211 }
212 return Err(ReadinessError::ConflictingVersion {
213 dependency,
214 version,
215 status,
216 });
217 }
218 current.evidence = ReadinessEvidence::Observed { version, status };
219 Ok(())
220 }
221
222 fn report(&self) -> ReadinessReport<K> {
223 ReadinessReport {
224 dependencies: self.dependencies.clone(),
225 }
226 }
227}
228
229impl<A, K, Route> BehaviorBase for Readiness<A, K, Route>
230where
231 A: Address,
232 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ReadinessReport<K>>>,
233{
234 type Base = Self;
235 fn base(&self) -> &Self {
236 self
237 }
238}
239
240impl<A, K, Route> behavior::Protocol for Readiness<A, K, Route>
241where
242 A: Address,
243 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ReadinessReport<K>>>,
244{
245 type Addr = A;
246 type Msg = ReadinessMessage<K, Route>;
247}
248
249impl<A, K, Route> Behavior for Readiness<A, K, Route>
250where
251 A: Address,
252 K: Clone + Eq,
253 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = ReadinessReport<K>>>,
254 Route::Sends: behavior::SendsFor<User<A, ReadinessMessage<K, Route>>>,
255{
256 type Protocol = Self;
257 type Event = User<A, behavior::BehaviorMessage<Self>>;
258 type Sends = Route::Sends;
259 type Ph = Never;
260 type Error = ReadinessError<K>;
261 type Birth = NoBirths;
262
263 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
264 match event.message {
265 ReadinessMessage::Observe {
266 dependency,
267 version,
268 status,
269 } => {
270 self.observe(dependency, version, status)?;
271 Ok(Actions::cont())
272 }
273 ReadinessMessage::Query { reply_to } => {
274 Ok(Actions::send(reply_to.deliver(self.report())))
275 }
276 }
277 }
278}
279
280#[cfg(test)]
281mod tests {
282 use super::*;
283 use crate::Activate as _;
284 use behavior::MailAddr;
285
286 struct Reply;
287 impl behavior::Protocol for Reply {
288 type Addr = MailAddr;
289 type Msg = ReadinessReport<u8>;
290 }
291
292 impl Behavior for Reply {
293 type Protocol = Self;
294 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
295 type Sends = Vec<Never>;
296 type Ph = Never;
297 type Error = Never;
298 type Birth = NoBirths;
299 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
300 Ok(Actions::cont())
301 }
302 }
303
304 type Subject = Readiness<MailAddr, u8, Recipient<Reply>>;
305
306 #[test]
307 fn all_dependencies_must_have_ready_evidence() {
308 let mut subject = (Subject::new([1, 1, 2])).initialize().unwrap().behavior;
309 let query = |subject: &mut crate::Active<Subject>| {
310 subject
311 .receive(
312 MailAddr(9),
313 ReadinessMessage::Query {
314 reply_to: Recipient::global(MailAddr(1)),
315 },
316 )
317 .unwrap()
318 };
319 let queried = query(&mut subject);
320 assert!(!queried.sends[0].message.ready());
321 let observed = subject
322 .receive(
323 MailAddr(9),
324 ReadinessMessage::Observe {
325 dependency: 1,
326 version: ObservationVersion(1),
327 status: ReadinessStatus::Ready,
328 },
329 )
330 .unwrap();
331 assert!(observed.sends.is_empty());
332 assert!(observed.creates.is_empty());
333 assert_eq!(observed.become_, behavior::Step::Continue);
334 let queried = query(&mut subject);
335 assert!(!queried.sends[0].message.ready());
336 let observed = subject
337 .receive(
338 MailAddr(9),
339 ReadinessMessage::Observe {
340 dependency: 2,
341 version: ObservationVersion(1),
342 status: ReadinessStatus::Ready,
343 },
344 )
345 .unwrap();
346 assert!(observed.sends.is_empty());
347 assert!(observed.creates.is_empty());
348 assert_eq!(observed.become_, behavior::Step::Continue);
349 let queried = query(&mut subject);
350 assert!(queried.sends[0].message.ready());
351 }
352
353 #[test]
354 fn stale_conflicting_and_unknown_evidence_are_atomic() {
355 let mut subject = (Subject::new([1])).initialize().unwrap().behavior;
356 let observed = subject
357 .receive(
358 MailAddr(9),
359 ReadinessMessage::Observe {
360 dependency: 1,
361 version: ObservationVersion(2),
362 status: ReadinessStatus::Ready,
363 },
364 )
365 .unwrap();
366 assert!(observed.sends.is_empty());
367 assert!(observed.creates.is_empty());
368 assert_eq!(observed.become_, behavior::Step::Continue);
369 let rejection = subject.receive(
370 MailAddr(9),
371 ReadinessMessage::Observe {
372 dependency: 1,
373 version: ObservationVersion(1),
374 status: ReadinessStatus::NotReady,
375 },
376 );
377 assert!(matches!(rejection, Err(ReadinessError::Stale { .. })));
378 let rejection = subject.receive(
379 MailAddr(9),
380 ReadinessMessage::Observe {
381 dependency: 1,
382 version: ObservationVersion(2),
383 status: ReadinessStatus::NotReady,
384 },
385 );
386 assert!(matches!(
387 rejection,
388 Err(ReadinessError::ConflictingVersion { .. })
389 ));
390 let rejection = subject.receive(
391 MailAddr(9),
392 ReadinessMessage::Observe {
393 dependency: 9,
394 version: ObservationVersion(1),
395 status: ReadinessStatus::Ready,
396 },
397 );
398 assert!(matches!(
399 rejection,
400 Err(ReadinessError::UnknownDependency {
401 dependency: 9,
402 observed: ObservationVersion(1),
403 status: ReadinessStatus::Ready,
404 })
405 ));
406 assert!(matches!(
407 subject.dependencies()[0].evidence,
408 ReadinessEvidence::Observed {
409 status: ReadinessStatus::Ready,
410 ..
411 }
412 ));
413 }
414
415 #[test]
416 fn empty_dependency_set_is_ready() {
417 let mut subject = (Subject::new([])).initialize().unwrap().behavior;
418 let report = subject
419 .receive(
420 MailAddr(9),
421 ReadinessMessage::Query {
422 reply_to: Recipient::global(MailAddr(1)),
423 },
424 )
425 .unwrap();
426 assert!(report.sends[0].message.ready());
427 }
428}