Skip to main content

soma_infra/
container_mutation_engine.rs

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/// Coordinates mutation and independent container-state verification.
11#[derive(Debug, Clone, Copy)]
12pub struct ContainerLifecycleEngine {
13    verification: MutationVerificationPolicy,
14}
15
16impl ContainerLifecycleEngine {
17    /// Creates an engine using the supplied verification policy.
18    #[must_use]
19    pub const fn new(verification: MutationVerificationPolicy) -> Self {
20        Self { verification }
21    }
22
23    /// Executes and verifies one lifecycle mutation.
24    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;