1#[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 ObservationVersion(pub u64);
13
14#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
16pub enum HealthStatus {
17 Healthy,
19 Degraded,
21 Unhealthy,
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct ComponentHealth<K> {
28 pub component: K,
30 pub version: ObservationVersion,
32 pub status: HealthStatus,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq)]
38pub enum ComponentHealthState<K> {
39 Present(ComponentHealth<K>),
41 Removed {
43 component: K,
45 version: ObservationVersion,
47 },
48}
49
50impl<K> ComponentHealthState<K> {
51 fn component(&self) -> &K {
52 match self {
53 Self::Present(health) => &health.component,
54 Self::Removed { component, .. } => component,
55 }
56 }
57
58 const fn version(&self) -> ObservationVersion {
59 match self {
60 Self::Present(health) => health.version,
61 Self::Removed { version, .. } => *version,
62 }
63 }
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct HealthReport<K> {
69 pub components: Vec<ComponentHealth<K>>,
71}
72
73impl<K> HealthReport<K> {
74 #[must_use]
76 pub fn overall(&self) -> HealthStatus {
77 self.components
78 .iter()
79 .map(|component| component.status)
80 .max()
81 .unwrap_or(HealthStatus::Healthy)
82 }
83}
84
85pub enum HealthMessage<K, Route> {
87 Observe {
89 component: K,
91 version: ObservationVersion,
93 status: HealthStatus,
95 },
96 Remove {
98 component: K,
100 version: ObservationVersion,
102 },
103 Query {
105 reply_to: Route,
107 },
108}
109
110#[derive(Debug, Error, Clone, PartialEq, Eq)]
112pub enum HealthError<K> {
113 #[error("health evidence is stale")]
115 Stale {
116 component: K,
118 observed: ObservationVersion,
120 current: ObservationVersion,
122 evidence: HealthEvidence,
124 },
125 #[error("health evidence contradicts the committed value at the same version")]
127 ConflictingVersion {
128 component: K,
130 version: ObservationVersion,
132 evidence: HealthEvidence,
134 },
135}
136
137#[derive(Debug, Clone, Copy, PartialEq, Eq)]
139pub enum HealthEvidence {
140 Present(HealthStatus),
142 Removed,
144}
145
146pub struct Health<A, K, Route>
159where
160 A: Address,
161 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
162{
163 components: Vec<ComponentHealthState<K>>,
164 address: core::marker::PhantomData<fn() -> (A, Route)>,
165}
166
167impl<A, K, Route> Health<A, K, Route>
168where
169 A: Address,
170 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
171{
172 #[must_use]
174 pub const fn new() -> Self {
175 Self {
176 components: Vec::new(),
177 address: core::marker::PhantomData,
178 }
179 }
180
181 #[must_use]
183 pub fn components(&self) -> &[ComponentHealthState<K>] {
184 &self.components
185 }
186}
187
188impl<A, K, Route> Default for Health<A, K, Route>
189where
190 A: Address,
191 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
192{
193 fn default() -> Self {
194 Self::new()
195 }
196}
197
198impl<A, K, Route> BehaviorBase for Health<A, K, Route>
199where
200 A: Address,
201 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
202{
203 type Base = Self;
204
205 fn base(&self) -> &Self {
206 self
207 }
208}
209
210impl<A, K, Route> Health<A, K, Route>
211where
212 A: Address,
213 K: Clone + Eq,
214 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
215{
216 fn commit(
217 &mut self,
218 component: K,
219 version: ObservationVersion,
220 evidence: HealthEvidence,
221 ) -> Result<(), HealthError<K>> {
222 let replacement = match evidence {
223 HealthEvidence::Present(status) => ComponentHealthState::Present(ComponentHealth {
224 component: component.clone(),
225 version,
226 status,
227 }),
228 HealthEvidence::Removed => ComponentHealthState::Removed {
229 component: component.clone(),
230 version,
231 },
232 };
233 let Some(index) = self
234 .components
235 .iter()
236 .position(|state| state.component() == &component)
237 else {
238 self.components.push(replacement);
239 return Ok(());
240 };
241 let current = self.components[index].version();
242 if version < current {
243 return Err(HealthError::Stale {
244 component,
245 observed: version,
246 current,
247 evidence,
248 });
249 }
250 if version == current {
251 if self.components[index] == replacement {
252 return Ok(());
253 }
254 return Err(HealthError::ConflictingVersion {
255 component,
256 version,
257 evidence,
258 });
259 }
260 self.components[index] = replacement;
261 Ok(())
262 }
263 fn report(&self) -> HealthReport<K> {
264 let components = self
265 .components
266 .iter()
267 .filter_map(|state| match state {
268 ComponentHealthState::Present(health) => Some(health.clone()),
269 ComponentHealthState::Removed { .. } => None,
270 })
271 .collect::<Vec<_>>();
272 HealthReport { components }
273 }
274}
275
276impl<A, K, Route> behavior::Protocol for Health<A, K, Route>
277where
278 A: Address,
279 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
280{
281 type Addr = A;
282 type Msg = HealthMessage<K, Route>;
283}
284
285impl<A, K, Route> Behavior for Health<A, K, Route>
286where
287 A: Address,
288 K: Clone + Eq,
289 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = HealthReport<K>>>,
290 Route::Sends: behavior::SendsFor<User<A, HealthMessage<K, Route>>>,
291{
292 type Protocol = Self;
293 type Event = User<A, behavior::BehaviorMessage<Self>>;
294 type Sends = Route::Sends;
295 type Ph = Never;
296 type Error = HealthError<K>;
297 type Birth = NoBirths;
298
299 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
300 match event.message {
301 HealthMessage::Observe {
302 component,
303 version,
304 status,
305 } => {
306 self.commit(component, version, HealthEvidence::Present(status))?;
307 Ok(Actions::cont())
308 }
309 HealthMessage::Remove { component, version } => {
310 self.commit(component, version, HealthEvidence::Removed)?;
311 Ok(Actions::cont())
312 }
313 HealthMessage::Query { reply_to } => Ok(Actions::send(reply_to.deliver(self.report()))),
314 }
315 }
316}
317
318#[cfg(test)]
319mod tests {
320 use super::*;
321 use crate::Activate as _;
322 use behavior::MailAddr;
323
324 struct Reply;
325
326 impl behavior::Protocol for Reply {
327 type Addr = MailAddr;
328 type Msg = HealthReport<u8>;
329 }
330
331 impl Behavior for Reply {
332 type Protocol = Self;
333 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
334 type Sends = Vec<Never>;
335 type Ph = Never;
336 type Error = Never;
337 type Birth = NoBirths;
338
339 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
340 Ok(Actions::cont())
341 }
342 }
343
344 type TestHealth = Health<MailAddr, u8, Recipient<Reply>>;
345
346 #[test]
347 fn stale_and_conflicting_evidence_preserve_committed_state() {
348 let mut health = (TestHealth::new()).initialize().unwrap().behavior;
349 let observed = health
350 .receive(
351 MailAddr(9),
352 HealthMessage::Observe {
353 component: 1,
354 version: ObservationVersion(3),
355 status: HealthStatus::Degraded,
356 },
357 )
358 .unwrap();
359 assert!(observed.sends.is_empty());
360 assert!(observed.creates.is_empty());
361 assert_eq!(observed.become_, behavior::Step::Continue);
362
363 let rejection = health.receive(
364 MailAddr(9),
365 HealthMessage::Observe {
366 component: 1,
367 version: ObservationVersion(2),
368 status: HealthStatus::Healthy,
369 },
370 );
371 assert!(matches!(
372 rejection,
373 Err(HealthError::Stale {
374 component: 1,
375 observed: ObservationVersion(2),
376 current: ObservationVersion(3),
377 evidence: HealthEvidence::Present(HealthStatus::Healthy),
378 })
379 ));
380 let rejection = health.receive(
381 MailAddr(9),
382 HealthMessage::Observe {
383 component: 1,
384 version: ObservationVersion(3),
385 status: HealthStatus::Unhealthy,
386 },
387 );
388 assert!(matches!(
389 rejection,
390 Err(HealthError::ConflictingVersion {
391 component: 1,
392 version: ObservationVersion(3),
393 evidence: HealthEvidence::Present(HealthStatus::Unhealthy),
394 })
395 ));
396 assert_eq!(
397 health.components(),
398 [ComponentHealthState::Present(ComponentHealth {
399 component: 1,
400 version: ObservationVersion(3),
401 status: HealthStatus::Degraded,
402 })]
403 );
404 }
405
406 #[test]
407 fn tombstone_rejects_resurrection_and_report_aggregates_worst_status() {
408 let reply = Recipient::<Reply>::global(MailAddr(8));
409 let mut health = (TestHealth::new()).initialize().unwrap().behavior;
410 for (component, status) in [(1, HealthStatus::Healthy), (2, HealthStatus::Unhealthy)] {
411 let observed = health
412 .receive(
413 MailAddr(9),
414 HealthMessage::Observe {
415 component,
416 version: ObservationVersion(1),
417 status,
418 },
419 )
420 .unwrap();
421 assert!(observed.sends.is_empty());
422 assert!(observed.creates.is_empty());
423 assert_eq!(observed.become_, behavior::Step::Continue);
424 }
425 let report = health
426 .receive(MailAddr(9), HealthMessage::Query { reply_to: reply })
427 .unwrap();
428 let sent = report
429 .sends
430 .into_iter()
431 .next()
432 .expect("one report delivery");
433 assert_eq!(sent.to, reply);
434 assert_eq!(sent.message.overall(), HealthStatus::Unhealthy);
435 assert_eq!(
436 sent.message.components,
437 vec![
438 ComponentHealth {
439 component: 1,
440 version: ObservationVersion(1),
441 status: HealthStatus::Healthy,
442 },
443 ComponentHealth {
444 component: 2,
445 version: ObservationVersion(1),
446 status: HealthStatus::Unhealthy,
447 },
448 ]
449 );
450
451 let removed = health
452 .receive(
453 MailAddr(9),
454 HealthMessage::Remove {
455 component: 2,
456 version: ObservationVersion(2),
457 },
458 )
459 .unwrap();
460 assert!(removed.sends.is_empty());
461 assert!(removed.creates.is_empty());
462 assert_eq!(removed.become_, behavior::Step::Continue);
463 let rejection = health.receive(
464 MailAddr(9),
465 HealthMessage::Observe {
466 component: 2,
467 version: ObservationVersion(1),
468 status: HealthStatus::Healthy,
469 },
470 );
471 assert!(matches!(rejection, Err(HealthError::Stale { .. })));
472 let after = health
473 .receive(MailAddr(9), HealthMessage::Query { reply_to: reply })
474 .unwrap();
475 let sent = after.sends.into_iter().next().expect("one report delivery");
476 assert_eq!(sent.to, reply);
477 assert_eq!(sent.message.overall(), HealthStatus::Healthy);
478 assert_eq!(
479 sent.message.components,
480 vec![ComponentHealth {
481 component: 1,
482 version: ObservationVersion(1),
483 status: HealthStatus::Healthy,
484 }]
485 );
486 }
487}