appcore_supervisor/
supervisor_lifecycle.rs1use super::*;
14
15impl Supervisor {
16 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 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 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 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 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}