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;