Skip to main content

appcore_supervisor/
supervisor_lifecycle.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: supervisor_lifecycle.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/07/24 13:18:47 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/07/24 13:18:47 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Non-blocking reconcile, dependency, watchdog, and shutdown lifecycle.
12
13use super::*;
14
15impl Supervisor {
16    /// Reconciles health and schedules restart work without executing it inline.
17    pub fn reconcile(&self, timestamp_ms: u64) -> SupervisorResult<()> {
18        self.inner.watchdog.record_reconcile_started(timestamp_ms);
19        self.apply_restart_completions(timestamp_ms)?;
20        for name in self.validate()? {
21            self.reconcile_service(&name, timestamp_ms)?;
22        }
23        self.dispatch_due_restarts(timestamp_ms)?;
24        let previous = self.inner.watchdog.state();
25        let sequence = self.inner.watchdog.record_reconcile_completed(timestamp_ms);
26        self.emit(
27            "supervisor",
28            SupervisorEventKind::SupervisorProgressed,
29            timestamp_ms,
30            sequence,
31            watchdog_states(previous, WatchdogState::Healthy),
32            "reconcile_completed",
33        );
34        Ok(())
35    }
36
37    /// Schedules one restart while enforcing its temporal budget.
38    pub fn restart(&self, name: &str, timestamp_ms: u64) -> SupervisorResult<()> {
39        let service = self.service(name)?;
40        self.require_dependencies(&service)?;
41        self.schedule_restart(&service, timestamp_ms)
42    }
43
44    /// Stops one enabled service without affecting the host process.
45    pub fn stop(&self, name: &str, timestamp_ms: u64) -> SupervisorResult<()> {
46        let service = self.service(name)?;
47        if !service.descriptor().activation().is_enabled() {
48            return Ok(());
49        }
50        let timeout = service.descriptor().restart_policy().shutdown_timeout;
51        match service.stop(timeout) {
52            Ok(()) => self.record_stopped(name, timestamp_ms),
53            Err(error) => {
54                self.record_stop_failure(&service, timestamp_ms)?;
55                Err(error)
56            }
57        }
58    }
59
60    /// Independently evaluates watchdog progress and emits transition events.
61    pub fn evaluate_watchdog(&self, timestamp_ms: u64) -> WatchdogSnapshot {
62        if let Some((previous, next)) = self.inner.watchdog.evaluate(timestamp_ms) {
63            let kind = match next {
64                WatchdogState::Stalled => SupervisorEventKind::SupervisorStalled,
65                WatchdogState::Healthy if previous == WatchdogState::Stalled => {
66                    SupervisorEventKind::SupervisorRecovered
67                }
68                _ => SupervisorEventKind::SupervisorProgressed,
69            };
70            self.emit(
71                "supervisor",
72                kind,
73                timestamp_ms,
74                self.inner.watchdog.reconcile_sequence(),
75                watchdog_states(previous, next),
76                "watchdog_evaluation",
77            );
78        }
79        self.watchdog_snapshot(timestamp_ms)
80    }
81
82    /// Stops the restart executor and all enabled services in reverse order.
83    pub fn shutdown(&self, timestamp_ms: u64) -> SupervisorResult<()> {
84        self.inner.watchdog.mark_stopping();
85        let _ = self
86            .inner
87            .restart_executor
88            .shutdown(Duration::from_secs(10));
89        let mut order = self.validate()?;
90        order.reverse();
91        let mut first_error = None;
92        for name in order {
93            let service = self.service(&name)?;
94            if !service.descriptor().activation().is_enabled() {
95                continue;
96            }
97            let timeout = service.descriptor().restart_policy().shutdown_timeout;
98            match service.stop(timeout) {
99                Ok(()) => self.record_stopped(&name, timestamp_ms)?,
100                Err(error) => {
101                    self.record_stop_failure(&service, timestamp_ms)?;
102                    if first_error.is_none() {
103                        first_error = Some(error);
104                    }
105                }
106            }
107        }
108        first_error.map_or(Ok(()), Err)
109    }
110
111    fn reconcile_service(&self, name: &str, timestamp_ms: u64) -> SupervisorResult<()> {
112        let service = self.service(name)?;
113        if !service.descriptor().activation().is_enabled() {
114            return Ok(());
115        }
116        if self.first_unavailable_dependency(&service)?.is_some() {
117            self.update_record(name, ServiceHealth::Degraded, service.runtime_state())?;
118            return Ok(());
119        }
120        let health = service.health();
121        let previous = self.update_record(name, health, service.runtime_state())?;
122        if health.is_failed() {
123            if previous != Some(ServiceHealth::Failed) {
124                self.emit(
125                    name,
126                    SupervisorEventKind::ServiceFailed,
127                    timestamp_ms,
128                    self.restart_attempt(name),
129                    states(previous, health),
130                    "health_failed",
131                );
132            }
133            if let Err(error) = self.schedule_restart(&service, timestamp_ms) {
134                if !matches!(error, SupervisorError::RestartBudgetExceeded(_)) {
135                    return Err(error);
136                }
137            }
138        } else if matches!(
139            previous,
140            Some(ServiceHealth::Failed | ServiceHealth::Degraded)
141        ) {
142            self.emit(
143                name,
144                SupervisorEventKind::ServiceRecovered,
145                timestamp_ms,
146                self.restart_attempt(name),
147                states(previous, health),
148                "health_recovered",
149            );
150        }
151        Ok(())
152    }
153
154    pub(super) fn require_dependencies(
155        &self,
156        service: &Arc<dyn ManagedService>,
157    ) -> SupervisorResult<()> {
158        if let Some(dependency) = self.first_unavailable_dependency(service)? {
159            return Err(SupervisorError::DependencyUnavailable {
160                service: service.descriptor().name().to_string(),
161                dependency,
162            });
163        }
164        Ok(())
165    }
166
167    fn first_unavailable_dependency(
168        &self,
169        service: &Arc<dyn ManagedService>,
170    ) -> SupervisorResult<Option<String>> {
171        for dependency in service.descriptor().dependencies() {
172            let dependency_service = match self.service(dependency.service_id()) {
173                Ok(service) => service,
174                Err(_) if dependency.requirement() == DependencyRequirement::Optional => continue,
175                Err(error) => return Err(error),
176            };
177            if !dependency
178                .requirement()
179                .accepts(dependency_service.health())
180            {
181                return Ok(Some(dependency.service_id().to_string()));
182            }
183        }
184        Ok(None)
185    }
186}