1use soma_fleet::HostRecord;
2use soma_ops::{MutationSendState, Timestamp, VerificationStatus};
3use tokio_util::sync::CancellationToken;
4
5use crate::{
6 ContainerLifecycleOutcome, ContainerLifecycleRequest, ContainerState, DockerMutationClient,
7 MutationFailure, MutationResult, MutationVerification, MutationVerificationPolicy,
8};
9
10#[derive(Debug, Clone, Copy)]
12pub struct ContainerLifecycleEngine {
13 verification: MutationVerificationPolicy,
14}
15
16impl ContainerLifecycleEngine {
17 #[must_use]
19 pub const fn new(verification: MutationVerificationPolicy) -> Self {
20 Self { verification }
21 }
22
23 pub async fn execute(
25 &self,
26 client: &dyn DockerMutationClient,
27 host: &HostRecord,
28 request: &ContainerLifecycleRequest,
29 cancellation: &CancellationToken,
30 ) -> MutationResult<ContainerLifecycleOutcome> {
31 ensure_admitted(request, cancellation)?;
32 let before = client
33 .inspect_container(host, request.container(), cancellation)
34 .await
35 .map_err(|error| MutationFailure::new(MutationSendState::NotSent, error))?
36 .state;
37
38 if request.action().already_satisfied(&before) {
39 return Ok(outcome(
40 host,
41 request,
42 false,
43 MutationSendState::NotSent,
44 before.clone(),
45 Some(before),
46 VerificationStatus::Verified,
47 "requested state was already satisfied",
48 ));
49 }
50
51 let receipt = client.mutate_container(host, request, cancellation).await?;
52 let verification = self.verify(client, host, request, cancellation).await;
53 Ok(outcome(
54 host,
55 request,
56 true,
57 receipt.send_state,
58 before,
59 verification.state,
60 verification.status,
61 verification.summary,
62 ))
63 }
64
65 async fn verify(
66 &self,
67 client: &dyn DockerMutationClient,
68 host: &HostRecord,
69 request: &ContainerLifecycleRequest,
70 cancellation: &CancellationToken,
71 ) -> VerificationObservation {
72 let mut last_state = None;
73 let mut last_error = None;
74 for attempt in 0..self.verification.attempts() {
75 if cancellation.is_cancelled() {
76 return VerificationObservation::inconclusive(
77 last_state,
78 "verification cancelled after mutation send",
79 );
80 }
81 if Timestamp::now() >= request.deadline() {
82 return VerificationObservation::inconclusive(
83 last_state,
84 "verification deadline expired after mutation send",
85 );
86 }
87 match client
88 .inspect_container(host, request.container(), cancellation)
89 .await
90 {
91 Ok(inspect) => {
92 last_state = Some(inspect.state);
93 if last_state
94 .as_ref()
95 .is_some_and(|state| request.action().verified(state))
96 {
97 return VerificationObservation {
98 state: last_state,
99 status: VerificationStatus::Verified,
100 summary: "runtime state matches the requested lifecycle state".into(),
101 };
102 }
103 }
104 Err(error) => last_error = Some(error.to_string()),
105 }
106 if attempt + 1 < self.verification.attempts() {
107 tokio::select! {
108 () = cancellation.cancelled() => {
109 return VerificationObservation::inconclusive(
110 last_state,
111 "verification cancelled after mutation send",
112 );
113 }
114 () = tokio::time::sleep(self.verification.interval()) => {}
115 }
116 }
117 }
118 let summary = last_error.map_or_else(
119 || "runtime state did not reach the requested lifecycle state".into(),
120 |error| format!("runtime state could not be verified: {error}"),
121 );
122 VerificationObservation {
123 state: last_state,
124 status: VerificationStatus::Failed,
125 summary,
126 }
127 }
128}
129
130impl Default for ContainerLifecycleEngine {
131 fn default() -> Self {
132 Self::new(MutationVerificationPolicy::default())
133 }
134}
135
136struct VerificationObservation {
137 state: Option<ContainerState>,
138 status: VerificationStatus,
139 summary: String,
140}
141
142impl VerificationObservation {
143 fn inconclusive(state: Option<ContainerState>, summary: &str) -> Self {
144 Self {
145 state,
146 status: VerificationStatus::Inconclusive,
147 summary: summary.into(),
148 }
149 }
150}
151
152fn ensure_admitted(
153 request: &ContainerLifecycleRequest,
154 cancellation: &CancellationToken,
155) -> MutationResult<()> {
156 if cancellation.is_cancelled() {
157 return Err(MutationFailure::new(
158 MutationSendState::NotSent,
159 soma_fleet::FleetError::Cancelled.into(),
160 ));
161 }
162 if Timestamp::now() >= request.deadline() {
163 return Err(MutationFailure::new(
164 MutationSendState::NotSent,
165 soma_fleet::FleetError::DeadlineExceeded.into(),
166 ));
167 }
168 Ok(())
169}
170
171#[allow(clippy::too_many_arguments)]
172fn outcome(
173 host: &HostRecord,
174 request: &ContainerLifecycleRequest,
175 changed: bool,
176 send_state: MutationSendState,
177 before: ContainerState,
178 after: Option<ContainerState>,
179 verification_status: VerificationStatus,
180 summary: impl Into<String>,
181) -> ContainerLifecycleOutcome {
182 ContainerLifecycleOutcome {
183 host: host.id().clone(),
184 topology_revision: host.revision().clone(),
185 container: request.container().to_owned(),
186 action: request.action(),
187 changed,
188 send_state,
189 before,
190 after,
191 verification_status,
192 verification: MutationVerification {
193 status: format!("{verification_status:?}").to_ascii_lowercase(),
194 summary: summary.into(),
195 },
196 }
197}
198
199#[cfg(test)]
200#[path = "container_mutation_engine_tests.rs"]
201mod tests;