ironflow_ops_docker/volumes/
mod.rs1use std::collections::HashMap;
4
5use async_trait::async_trait;
6use bollard::Docker;
7use bollard::models::VolumeCreateRequest;
8use bollard::query_parameters::{ListVolumesOptions, PruneVolumesOptions, RemoveVolumeOptions};
9use ironflow_core::error::OperationError;
10use ironflow_core::operation::{Operation, OperationContext, TypedOperation};
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13
14use crate::containers::DockerRef;
15use crate::helpers::{docker_error, to_value};
16
17#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct VolumeCreateOutput {
24 pub name: String,
26 pub mountpoint: String,
28}
29
30pub struct VolumeCreate {
44 docker: Docker,
45 name: String,
46 driver: Option<String>,
47 labels: HashMap<String, String>,
48}
49
50impl VolumeCreate {
51 pub fn new(client: impl Into<DockerRef>, name: impl Into<String>) -> Self {
53 Self {
54 docker: client.into().0,
55 name: name.into(),
56 driver: None,
57 labels: HashMap::new(),
58 }
59 }
60
61 pub fn driver(mut self, driver: impl Into<String>) -> Self {
63 self.driver = Some(driver.into());
64 self
65 }
66
67 pub fn label(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
69 self.labels.insert(key.into(), value.into());
70 self
71 }
72
73 pub async fn run(&self, _ctx: &OperationContext) -> Result<VolumeCreateOutput, OperationError> {
79 let request = VolumeCreateRequest {
80 name: Some(self.name.clone()),
81 driver: Some(self.driver.clone().unwrap_or_else(|| "local".to_string())),
82 labels: Some(self.labels.clone()),
83 ..Default::default()
84 };
85 let response = self
86 .docker
87 .create_volume(request)
88 .await
89 .map_err(docker_error)?;
90 Ok(VolumeCreateOutput {
91 name: response.name,
92 mountpoint: response.mountpoint,
93 })
94 }
95}
96
97#[async_trait]
98impl Operation for VolumeCreate {
99 fn kind(&self) -> &str {
100 "docker"
101 }
102
103 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
104 to_value(&self.run(ctx).await?)
105 }
106
107 fn input(&self) -> Option<Value> {
108 Some(serde_json::json!({
109 "operation": "volume_create",
110 "name": self.name,
111 }))
112 }
113}
114
115impl TypedOperation for VolumeCreate {
116 type Output = VolumeCreateOutput;
117}
118
119#[derive(Debug, Clone, Serialize, Deserialize)]
125pub struct VolumeInspectOutput {
126 pub name: String,
128 pub mountpoint: String,
130 pub driver: String,
132}
133
134pub struct VolumeInspect {
148 docker: Docker,
149 name: String,
150}
151
152impl VolumeInspect {
153 pub fn new(client: impl Into<DockerRef>, name: impl Into<String>) -> Self {
155 Self {
156 docker: client.into().0,
157 name: name.into(),
158 }
159 }
160
161 pub async fn run(
167 &self,
168 _ctx: &OperationContext,
169 ) -> Result<VolumeInspectOutput, OperationError> {
170 let response = self
171 .docker
172 .inspect_volume(&self.name)
173 .await
174 .map_err(docker_error)?;
175 Ok(VolumeInspectOutput {
176 name: response.name,
177 mountpoint: response.mountpoint,
178 driver: response.driver,
179 })
180 }
181}
182
183#[async_trait]
184impl Operation for VolumeInspect {
185 fn kind(&self) -> &str {
186 "docker"
187 }
188
189 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
190 to_value(&self.run(ctx).await?)
191 }
192
193 fn input(&self) -> Option<Value> {
194 Some(serde_json::json!({
195 "operation": "volume_inspect",
196 "name": self.name,
197 }))
198 }
199}
200
201impl TypedOperation for VolumeInspect {
202 type Output = VolumeInspectOutput;
203}
204
205#[derive(Debug, Clone, Serialize, Deserialize)]
211pub struct VolumeListEntry {
212 pub name: String,
214 pub driver: String,
216 pub mountpoint: String,
218}
219
220#[derive(Debug, Clone, Serialize, Deserialize)]
222pub struct VolumeListOutput {
223 pub volumes: Vec<VolumeListEntry>,
225}
226
227pub struct VolumeList {
241 docker: Docker,
242 filters: HashMap<String, Vec<String>>,
243}
244
245impl VolumeList {
246 pub fn new(client: impl Into<DockerRef>) -> Self {
248 Self {
249 docker: client.into().0,
250 filters: HashMap::new(),
251 }
252 }
253
254 pub fn filter(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
256 self.filters
257 .entry(key.into())
258 .or_default()
259 .push(value.into());
260 self
261 }
262
263 pub async fn run(&self, _ctx: &OperationContext) -> Result<VolumeListOutput, OperationError> {
270 let options = ListVolumesOptions {
271 filters: Some(self.filters.clone()),
272 };
273 let response = self
274 .docker
275 .list_volumes(Some(options))
276 .await
277 .map_err(docker_error)?;
278 let volumes = response
279 .volumes
280 .unwrap_or_default()
281 .into_iter()
282 .map(|v| VolumeListEntry {
283 name: v.name,
284 driver: v.driver,
285 mountpoint: v.mountpoint,
286 })
287 .collect();
288 Ok(VolumeListOutput { volumes })
289 }
290}
291
292#[async_trait]
293impl Operation for VolumeList {
294 fn kind(&self) -> &str {
295 "docker"
296 }
297
298 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
299 to_value(&self.run(ctx).await?)
300 }
301
302 fn input(&self) -> Option<Value> {
303 Some(serde_json::json!({
304 "operation": "volume_list",
305 }))
306 }
307}
308
309impl TypedOperation for VolumeList {
310 type Output = VolumeListOutput;
311}
312
313#[derive(Debug, Clone, Serialize, Deserialize)]
319pub struct VolumeRemoveOutput {
320 pub name: String,
322}
323
324pub struct VolumeRemove {
338 docker: Docker,
339 name: String,
340 force: bool,
341}
342
343impl VolumeRemove {
344 pub fn new(client: impl Into<DockerRef>, name: impl Into<String>) -> Self {
346 Self {
347 docker: client.into().0,
348 name: name.into(),
349 force: false,
350 }
351 }
352
353 pub fn force(mut self) -> Self {
355 self.force = true;
356 self
357 }
358
359 pub async fn run(&self, _ctx: &OperationContext) -> Result<VolumeRemoveOutput, OperationError> {
365 self.docker
366 .remove_volume(&self.name, Some(RemoveVolumeOptions { force: self.force }))
367 .await
368 .map_err(docker_error)?;
369 Ok(VolumeRemoveOutput {
370 name: self.name.clone(),
371 })
372 }
373}
374
375#[async_trait]
376impl Operation for VolumeRemove {
377 fn kind(&self) -> &str {
378 "docker"
379 }
380
381 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
382 to_value(&self.run(ctx).await?)
383 }
384
385 fn input(&self) -> Option<Value> {
386 Some(serde_json::json!({
387 "operation": "volume_remove",
388 "name": self.name,
389 }))
390 }
391}
392
393impl TypedOperation for VolumeRemove {
394 type Output = VolumeRemoveOutput;
395}
396
397#[derive(Debug, Clone, Serialize, Deserialize)]
403pub struct VolumePruneOutput {
404 pub volumes_deleted: Vec<String>,
406 pub space_reclaimed: u64,
408}
409
410pub struct VolumePrune {
424 docker: Docker,
425 filters: HashMap<String, Vec<String>>,
426}
427
428impl VolumePrune {
429 pub fn new(client: impl Into<DockerRef>) -> Self {
431 Self {
432 docker: client.into().0,
433 filters: HashMap::new(),
434 }
435 }
436
437 pub fn filter(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
439 self.filters
440 .entry(key.into())
441 .or_default()
442 .push(value.into());
443 self
444 }
445
446 pub async fn run(&self, _ctx: &OperationContext) -> Result<VolumePruneOutput, OperationError> {
453 let options = PruneVolumesOptions {
454 filters: Some(self.filters.clone()),
455 };
456 let response = self
457 .docker
458 .prune_volumes(Some(options))
459 .await
460 .map_err(docker_error)?;
461 Ok(VolumePruneOutput {
462 volumes_deleted: response.volumes_deleted.unwrap_or_default(),
463 space_reclaimed: response.space_reclaimed.unwrap_or(0) as u64,
464 })
465 }
466}
467
468#[async_trait]
469impl Operation for VolumePrune {
470 fn kind(&self) -> &str {
471 "docker"
472 }
473
474 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
475 to_value(&self.run(ctx).await?)
476 }
477
478 fn input(&self) -> Option<Value> {
479 Some(serde_json::json!({
480 "operation": "volume_prune",
481 }))
482 }
483}
484
485impl TypedOperation for VolumePrune {
486 type Output = VolumePruneOutput;
487}