1use async_trait::async_trait;
2use serde::{Deserialize, Serialize};
3use soma_fleet::{HostId, HostRecord, TopologyRevision};
4use soma_ops::{MutationSendState, Timestamp, VerificationStatus};
5use tokio_util::sync::CancellationToken;
6
7use crate::{
8 ComposeInspector, ComposeProjectRef, ComposeStatus, MutationFailure, MutationResult,
9 MutationVerification, MutationVerificationPolicy,
10};
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
14#[serde(rename_all = "snake_case")]
15pub enum ComposeMutationAction {
16 Up,
18 Restart,
20}
21
22impl ComposeMutationAction {
23 #[must_use]
25 pub const fn operation_name(self) -> &'static str {
26 match self {
27 Self::Up => "compose.up",
28 Self::Restart => "compose.restart",
29 }
30 }
31
32 #[must_use]
34 pub const fn action_label(self) -> &'static str {
35 match self {
36 Self::Up => "up",
37 Self::Restart => "restart",
38 }
39 }
40}
41
42#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
44pub struct ComposeMutationRequest {
45 project: ComposeProjectRef,
46 action: ComposeMutationAction,
47 deadline: Timestamp,
48}
49
50impl ComposeMutationRequest {
51 #[must_use]
53 pub const fn new(
54 project: ComposeProjectRef,
55 action: ComposeMutationAction,
56 deadline: Timestamp,
57 ) -> Self {
58 Self {
59 project,
60 action,
61 deadline,
62 }
63 }
64
65 #[must_use]
67 pub const fn project(&self) -> &ComposeProjectRef {
68 &self.project
69 }
70
71 #[must_use]
73 pub const fn action(&self) -> ComposeMutationAction {
74 self.action
75 }
76
77 #[must_use]
79 pub const fn deadline(&self) -> Timestamp {
80 self.deadline
81 }
82}
83
84#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
86pub struct ComposeMutationReceipt {
87 pub host: HostId,
89 pub topology_revision: TopologyRevision,
91 pub project: String,
93 pub action: ComposeMutationAction,
95 pub send_state: MutationSendState,
97}
98
99#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101pub struct ComposeMutationOutcome {
102 pub host: HostId,
104 pub topology_revision: TopologyRevision,
106 pub project: String,
108 pub action: ComposeMutationAction,
110 pub send_state: MutationSendState,
112 pub before: Option<ComposeStatus>,
114 pub after: Option<ComposeStatus>,
116 pub verification_status: VerificationStatus,
118 pub verification: MutationVerification,
120}
121
122#[async_trait]
124pub trait ComposeMutator: Send + Sync {
125 async fn mutate_compose(
127 &self,
128 host: &HostRecord,
129 request: &ComposeMutationRequest,
130 cancellation: &CancellationToken,
131 ) -> MutationResult<ComposeMutationReceipt>;
132}
133
134pub trait ComposeMutationClient: ComposeInspector + ComposeMutator {}
136
137impl<T> ComposeMutationClient for T where T: ComposeInspector + ComposeMutator {}
138
139#[derive(Debug, Clone, Copy)]
141pub struct ComposeMutationEngine {
142 verification: MutationVerificationPolicy,
143}
144
145impl ComposeMutationEngine {
146 #[must_use]
148 pub const fn new(verification: MutationVerificationPolicy) -> Self {
149 Self { verification }
150 }
151
152 pub async fn execute(
154 &self,
155 client: &dyn ComposeMutationClient,
156 host: &HostRecord,
157 request: &ComposeMutationRequest,
158 cancellation: &CancellationToken,
159 ) -> MutationResult<ComposeMutationOutcome> {
160 ensure_admitted(request, cancellation)?;
161 let before = client
162 .status(
163 host,
164 request.project(),
165 None,
166 request.deadline(),
167 cancellation,
168 )
169 .await
170 .ok();
171 let receipt = client.mutate_compose(host, request, cancellation).await?;
172 let verification = self.verify(client, host, request, cancellation).await;
173 Ok(ComposeMutationOutcome {
174 host: host.id().clone(),
175 topology_revision: host.revision().clone(),
176 project: request.project().name().to_owned(),
177 action: request.action(),
178 send_state: receipt.send_state,
179 before,
180 after: verification.status,
181 verification_status: verification.verification_status,
182 verification: MutationVerification {
183 status: verification_status_text(verification.verification_status),
184 summary: verification.summary,
185 },
186 })
187 }
188
189 async fn verify(
190 &self,
191 client: &dyn ComposeMutationClient,
192 host: &HostRecord,
193 request: &ComposeMutationRequest,
194 cancellation: &CancellationToken,
195 ) -> ComposeVerificationObservation {
196 let mut last_status = None;
197 let mut last_error = None;
198 for attempt in 0..self.verification.attempts() {
199 if cancellation.is_cancelled() {
200 return ComposeVerificationObservation::inconclusive(
201 last_status,
202 "Compose verification cancelled after mutation send",
203 );
204 }
205 if Timestamp::now() >= request.deadline() {
206 return ComposeVerificationObservation::inconclusive(
207 last_status,
208 "Compose verification deadline expired after mutation send",
209 );
210 }
211 match client
212 .status(
213 host,
214 request.project(),
215 None,
216 request.deadline(),
217 cancellation,
218 )
219 .await
220 {
221 Ok(status) => {
222 if compose_status_running(&status) {
223 return ComposeVerificationObservation {
224 status: Some(status),
225 verification_status: VerificationStatus::Verified,
226 summary: "all reported Compose services are running".into(),
227 };
228 }
229 last_status = Some(status);
230 }
231 Err(error) => last_error = Some(error.to_string()),
232 }
233 if attempt + 1 < self.verification.attempts() {
234 tokio::select! {
235 () = cancellation.cancelled() => {
236 return ComposeVerificationObservation::inconclusive(
237 last_status,
238 "Compose verification cancelled after mutation send",
239 );
240 }
241 () = tokio::time::sleep(self.verification.interval()) => {}
242 }
243 }
244 }
245 let summary = last_error.map_or_else(
246 || "one or more Compose services did not reach running state".into(),
247 |error| format!("Compose status could not be verified: {error}"),
248 );
249 ComposeVerificationObservation {
250 status: last_status,
251 verification_status: VerificationStatus::Failed,
252 summary,
253 }
254 }
255}
256
257impl Default for ComposeMutationEngine {
258 fn default() -> Self {
259 Self::new(MutationVerificationPolicy::default())
260 }
261}
262
263struct ComposeVerificationObservation {
264 status: Option<ComposeStatus>,
265 verification_status: VerificationStatus,
266 summary: String,
267}
268
269impl ComposeVerificationObservation {
270 fn inconclusive(status: Option<ComposeStatus>, summary: &str) -> Self {
271 Self {
272 status,
273 verification_status: VerificationStatus::Inconclusive,
274 summary: summary.into(),
275 }
276 }
277}
278
279fn ensure_admitted(
280 request: &ComposeMutationRequest,
281 cancellation: &CancellationToken,
282) -> MutationResult<()> {
283 if cancellation.is_cancelled() {
284 return Err(MutationFailure::new(
285 MutationSendState::NotSent,
286 soma_fleet::FleetError::Cancelled.into(),
287 ));
288 }
289 if Timestamp::now() >= request.deadline() {
290 return Err(MutationFailure::new(
291 MutationSendState::NotSent,
292 soma_fleet::FleetError::DeadlineExceeded.into(),
293 ));
294 }
295 Ok(())
296}
297
298fn compose_status_running(status: &ComposeStatus) -> bool {
299 !status.services.is_empty()
300 && status.services.iter().all(|service| {
301 service
302 .state
303 .as_deref()
304 .is_some_and(|state| state.eq_ignore_ascii_case("running"))
305 && !service.health.as_deref().is_some_and(|health| {
306 matches!(
307 health.to_ascii_lowercase().as_str(),
308 "unhealthy" | "starting"
309 )
310 })
311 && service.exit_code.unwrap_or(0) == 0
312 })
313}
314
315fn verification_status_text(status: VerificationStatus) -> String {
316 format!("{status:?}").to_ascii_lowercase()
317}
318
319#[cfg(test)]
320#[path = "compose_mutation_tests.rs"]
321mod tests;