Skip to main content

behavior_actors/operations/
health.rs

1//! Versioned component-health aggregation.
2
3#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10/// Monotonic version of one component's health evidence.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
12pub struct ObservationVersion(pub u64);
13
14/// Closed health classification ordered from best to worst.
15#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
16pub enum HealthStatus {
17    /// The component is fully available under its declared policy.
18    Healthy,
19    /// The component remains usable with a declared impairment.
20    Degraded,
21    /// The component is unavailable or must not receive work.
22    Unhealthy,
23}
24
25/// One present component in a [`HealthReport`].
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct ComponentHealth<K> {
28    /// Application-defined component identity.
29    pub component: K,
30    /// Latest committed evidence version.
31    pub version: ObservationVersion,
32    /// Latest committed classification.
33    pub status: HealthStatus,
34}
35
36/// Complete retained state for one known component.
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub enum ComponentHealthState<K> {
39    /// Latest evidence describes a present component.
40    Present(ComponentHealth<K>),
41    /// The component was explicitly removed at this version.
42    Removed {
43        /// Application-defined component identity.
44        component: K,
45        /// Removal version retained to reject stale resurrection.
46        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/// Factual point-in-time result returned by [`Health`].
67#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct HealthReport<K> {
69    /// Present components in first-observation order; tombstones are omitted.
70    pub components: Vec<ComponentHealth<K>>,
71}
72
73impl<K> HealthReport<K> {
74    /// Worst status among present components, or `Healthy` when none exist.
75    #[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
85/// Commands accepted by [`Health`].
86pub enum HealthMessage<K, Route> {
87    /// Commit versioned evidence for a present component.
88    Observe {
89        /// Component identity.
90        component: K,
91        /// Evidence version within that component's stream.
92        version: ObservationVersion,
93        /// New factual classification.
94        status: HealthStatus,
95    },
96    /// Commit a versioned tombstone.
97    Remove {
98        /// Component identity.
99        component: K,
100        /// Removal version within that component's stream.
101        version: ObservationVersion,
102    },
103    /// Return a point-in-time report to a typed recipient.
104    Query {
105        /// Recipient whose protocol accepts [`HealthReport<K>`].
106        reply_to: Route,
107    },
108}
109
110/// Rejected versioned health evidence.
111#[derive(Debug, Error, Clone, PartialEq, Eq)]
112pub enum HealthError<K> {
113    /// Evidence predates the latest committed version.
114    #[error("health evidence is stale")]
115    Stale {
116        /// Component whose evidence was rejected.
117        component: K,
118        /// Rejected version.
119        observed: ObservationVersion,
120        /// Latest committed version.
121        current: ObservationVersion,
122        /// Exact rejected evidence.
123        evidence: HealthEvidence,
124    },
125    /// Evidence reuses a committed version with different meaning.
126    #[error("health evidence contradicts the committed value at the same version")]
127    ConflictingVersion {
128        /// Component whose evidence was rejected.
129        component: K,
130        /// Reused version.
131        version: ObservationVersion,
132        /// Exact rejected evidence.
133        evidence: HealthEvidence,
134    },
135}
136
137/// Exact meaning of one submitted health fact.
138#[derive(Debug, Clone, Copy, PartialEq, Eq)]
139pub enum HealthEvidence {
140    /// The component is present with this classification.
141    Present(HealthStatus),
142    /// The component is removed at the submitted version.
143    Removed,
144}
145
146/// Versioned typed health aggregation behavior.
147///
148/// State is an insertion-ordered product of [`ComponentHealthState`] values.
149/// A greater version replaces the prior state. An identical observation is
150/// idempotent. A lower version is [`HealthError::Stale`], and a same-version
151/// contradiction is [`HealthError::ConflictingVersion`]; both preserve the
152/// complete prior state. Removal retains a tombstone, so stale evidence cannot
153/// resurrect a component. Query emits one named typed delivery and never
154/// mutates state. Initialization is empty, the behavior never terminates by
155/// policy, and it requires only ordinary typed delivery interpretation.
156/// Versioning, aggregate ordering, and empty-set health are Bombay policy, not
157/// actor-model laws. No method has a semantic panic condition.
158pub 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    /// Construct an empty health definition.
173    #[must_use]
174    pub const fn new() -> Self {
175        Self {
176            components: Vec::new(),
177            address: core::marker::PhantomData,
178        }
179    }
180
181    /// Borrow every retained present or removed component state.
182    #[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}