Skip to main content

soma_infra/
bollard_cleanup.rs

1use std::future::Future;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use bollard::query_parameters::{
6    PruneBuildOptions, PruneContainersOptions, PruneImagesOptions, PruneNetworksOptions,
7    PruneVolumesOptions, RemoveImageOptionsBuilder,
8};
9use soma_fleet::HostRecord;
10use soma_ops::{MutationSendState, Timestamp};
11use tokio_util::sync::CancellationToken;
12
13use crate::{
14    BollardReadClient, DockerCleanupMutator, DockerPruneReceipt, DockerPruneRequest,
15    DockerPruneScopeReceipt, DockerPruneTarget, ImageRemovalReceipt, ImageRemovalRequest,
16    InfraError, MutationFailure, MutationResult,
17};
18
19#[async_trait]
20impl DockerCleanupMutator for BollardReadClient {
21    async fn remove_image(
22        &self,
23        host: &HostRecord,
24        request: &ImageRemovalRequest,
25        cancellation: &CancellationToken,
26    ) -> MutationResult<ImageRemovalReceipt> {
27        self.validate_host(host).map_err(not_sent)?;
28        ensure_admitted(request.force, request.deadline, cancellation)?;
29        let options = RemoveImageOptionsBuilder::default()
30            .force(request.force)
31            .build();
32        let rows = await_send(
33            request.deadline,
34            cancellation,
35            self.docker()
36                .remove_image(&request.fingerprint.reference, Some(options), None),
37        )
38        .await?;
39        let mut deleted = rows
40            .iter()
41            .filter_map(|row| row.deleted.clone())
42            .collect::<Vec<_>>();
43        let mut untagged = rows
44            .iter()
45            .filter_map(|row| row.untagged.clone())
46            .collect::<Vec<_>>();
47        deleted.sort();
48        deleted.dedup();
49        untagged.sort();
50        untagged.dedup();
51        Ok(ImageRemovalReceipt {
52            send_state: MutationSendState::Sent,
53            deleted,
54            untagged,
55        })
56    }
57
58    async fn prune(
59        &self,
60        host: &HostRecord,
61        request: &DockerPruneRequest,
62        cancellation: &CancellationToken,
63    ) -> MutationResult<DockerPruneReceipt> {
64        self.validate_host(host).map_err(not_sent)?;
65        ensure_admitted(request.force, request.deadline, cancellation)?;
66        let mut scopes = Vec::new();
67        for target in request.fingerprint.target.expanded() {
68            let scope = self
69                .prune_scope(*target, request.deadline, cancellation)
70                .await
71                .map_err(|failure| {
72                    let completed = scopes
73                        .iter()
74                        .map(|scope: &DockerPruneScopeReceipt| scope.target.as_str())
75                        .collect::<Vec<_>>()
76                        .join(",");
77                    MutationFailure::new(
78                        failure.send_state(),
79                        InfraError::Docker(format!(
80                            "prune scope {} failed after completed scopes [{completed}]: {}",
81                            target.as_str(),
82                            failure.error()
83                        )),
84                    )
85                })?;
86            scopes.push(scope);
87        }
88        Ok(DockerPruneReceipt {
89            send_state: MutationSendState::Sent,
90            scopes,
91        })
92    }
93}
94
95impl BollardReadClient {
96    async fn prune_scope(
97        &self,
98        target: DockerPruneTarget,
99        deadline: Timestamp,
100        cancellation: &CancellationToken,
101    ) -> MutationResult<DockerPruneScopeReceipt> {
102        match target {
103            DockerPruneTarget::Containers => {
104                let response = await_send(
105                    deadline,
106                    cancellation,
107                    self.docker()
108                        .prune_containers(None::<PruneContainersOptions>),
109                )
110                .await?;
111                Ok(scope(
112                    target,
113                    response.containers_deleted.unwrap_or_default(),
114                    response.space_reclaimed,
115                ))
116            }
117            DockerPruneTarget::Images => {
118                let response = await_send(
119                    deadline,
120                    cancellation,
121                    self.docker().prune_images(None::<PruneImagesOptions>),
122                )
123                .await?;
124                let deleted = response
125                    .images_deleted
126                    .unwrap_or_default()
127                    .into_iter()
128                    .filter_map(|row| row.deleted)
129                    .collect();
130                Ok(scope(target, deleted, response.space_reclaimed))
131            }
132            DockerPruneTarget::Volumes => {
133                let response = await_send(
134                    deadline,
135                    cancellation,
136                    self.docker().prune_volumes(None::<PruneVolumesOptions>),
137                )
138                .await?;
139                Ok(scope(
140                    target,
141                    response.volumes_deleted.unwrap_or_default(),
142                    response.space_reclaimed,
143                ))
144            }
145            DockerPruneTarget::Networks => {
146                let response = await_send(
147                    deadline,
148                    cancellation,
149                    self.docker().prune_networks(None::<PruneNetworksOptions>),
150                )
151                .await?;
152                Ok(scope(
153                    target,
154                    response.networks_deleted.unwrap_or_default(),
155                    None,
156                ))
157            }
158            DockerPruneTarget::BuildCache => {
159                let response = await_send(
160                    deadline,
161                    cancellation,
162                    self.docker().prune_build(None::<PruneBuildOptions>),
163                )
164                .await?;
165                Ok(scope(
166                    target,
167                    response.caches_deleted.unwrap_or_default(),
168                    response.space_reclaimed,
169                ))
170            }
171            DockerPruneTarget::All => unreachable!("expanded before execution"),
172        }
173    }
174}
175
176fn scope(
177    target: DockerPruneTarget,
178    mut deleted: Vec<String>,
179    reclaimed: Option<i64>,
180) -> DockerPruneScopeReceipt {
181    deleted.sort();
182    deleted.dedup();
183    DockerPruneScopeReceipt {
184        target,
185        deleted,
186        space_reclaimed: reclaimed.unwrap_or_default().max(0) as u64,
187    }
188}
189
190fn ensure_admitted(
191    force: bool,
192    deadline: Timestamp,
193    cancellation: &CancellationToken,
194) -> MutationResult<()> {
195    if !force {
196        return Err(not_sent(InfraError::InvalidRequest {
197            domain: "docker-cleanup",
198            message: "force=true is required".into(),
199        }));
200    }
201    if cancellation.is_cancelled() {
202        return Err(not_sent(soma_fleet::FleetError::Cancelled.into()));
203    }
204    if deadline <= Timestamp::now() {
205        return Err(not_sent(soma_fleet::FleetError::DeadlineExceeded.into()));
206    }
207    Ok(())
208}
209
210async fn await_send<T, F>(
211    deadline: Timestamp,
212    cancellation: &CancellationToken,
213    future: F,
214) -> MutationResult<T>
215where
216    F: Future<Output = Result<T, bollard::errors::Error>>,
217{
218    let remaining = deadline
219        .unix_millis()
220        .saturating_sub(Timestamp::now().unix_millis());
221    if remaining <= 0 {
222        return Err(not_sent(soma_fleet::FleetError::DeadlineExceeded.into()));
223    }
224    tokio::select! {
225        () = cancellation.cancelled() => Err(MutationFailure::new(
226            MutationSendState::Unknown,
227            soma_fleet::FleetError::Cancelled.into(),
228        )),
229        () = tokio::time::sleep(Duration::from_millis(remaining as u64)) => Err(
230            MutationFailure::new(
231                MutationSendState::Unknown,
232                soma_fleet::FleetError::DeadlineExceeded.into(),
233            )
234        ),
235        result = future => result.map_err(|error| MutationFailure::new(
236            MutationSendState::Unknown,
237            InfraError::Docker(error.to_string()),
238        )),
239    }
240}
241
242fn not_sent(error: InfraError) -> MutationFailure {
243    MutationFailure::new(MutationSendState::NotSent, error)
244}
245
246#[cfg(test)]
247#[path = "bollard_cleanup_tests.rs"]
248mod tests;