Skip to main content

ironflow_ops_docker/containers/
cleanup.rs

1//! Container cleanup operations: stats, changes, prune.
2
3use std::collections::HashMap;
4
5use async_trait::async_trait;
6use bollard::Docker;
7use bollard::query_parameters::{PruneContainersOptions, StatsOptions};
8use ironflow_core::error::OperationError;
9use ironflow_core::operation::{Operation, OperationContext, TypedOperation};
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12use tokio_stream::StreamExt;
13
14use crate::containers::DockerRef;
15use crate::helpers::{docker_error, to_value};
16
17// ---------------------------------------------------------------------------
18// ContainerStats
19// ---------------------------------------------------------------------------
20
21/// Output of a single stats snapshot.
22#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct ContainerStatsOutput {
24    /// The stats data as JSON.
25    pub data: Value,
26}
27
28/// Get a single stats snapshot from a container.
29///
30/// # Examples
31///
32/// ```no_run
33/// use ironflow_ops_docker::containers::ContainerStats;
34/// use ironflow_ops_docker::DockerClient;
35/// use ironflow_core::operation::Operation;
36///
37/// let client = DockerClient::connect_local().unwrap();
38/// let op = ContainerStats::new(&client, "my-container");
39/// assert_eq!(op.kind(), "docker");
40/// ```
41pub struct ContainerStats {
42    docker: Docker,
43    container: String,
44}
45
46impl ContainerStats {
47    /// Create a new container-stats operation (single snapshot).
48    pub fn new(client: impl Into<DockerRef>, container: impl Into<String>) -> Self {
49        Self {
50            docker: client.into().0,
51            container: container.into(),
52        }
53    }
54
55    /// Execute and return a typed result.
56    ///
57    /// # Errors
58    ///
59    /// Returns [`OperationError::External`] if the container does not exist.
60    pub async fn run(
61        &self,
62        _ctx: &OperationContext,
63    ) -> Result<ContainerStatsOutput, OperationError> {
64        let options = StatsOptions {
65            stream: false,
66            one_shot: true,
67        };
68        let mut stream = self.docker.stats(&self.container, Some(options));
69        let stats = stream
70            .next()
71            .await
72            .ok_or_else(|| OperationError::External {
73                origin: "docker".to_string(),
74                message: "no stats returned".to_string(),
75            })?
76            .map_err(docker_error)?;
77        let data = serde_json::to_value(&stats).map_err(|e| OperationError::External {
78            origin: "docker".to_string(),
79            message: e.to_string(),
80        })?;
81        Ok(ContainerStatsOutput { data })
82    }
83}
84
85#[async_trait]
86impl Operation for ContainerStats {
87    fn kind(&self) -> &str {
88        "docker"
89    }
90
91    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
92        to_value(&self.run(ctx).await?)
93    }
94
95    fn input(&self) -> Option<Value> {
96        Some(serde_json::json!({
97            "operation": "container_stats",
98            "container": self.container,
99        }))
100    }
101}
102
103impl TypedOperation for ContainerStats {
104    type Output = ContainerStatsOutput;
105}
106
107// ---------------------------------------------------------------------------
108// ContainerChanges
109// ---------------------------------------------------------------------------
110
111/// A single filesystem change in a container.
112#[derive(Debug, Clone, Serialize, Deserialize)]
113pub struct ContainerChange {
114    /// The file path.
115    pub path: String,
116    /// The kind of change (0=Modified, 1=Added, 2=Deleted).
117    pub kind: i32,
118}
119
120/// Output of listing container filesystem changes.
121#[derive(Debug, Clone, Serialize, Deserialize)]
122pub struct ContainerChangesOutput {
123    /// The filesystem changes.
124    pub changes: Vec<ContainerChange>,
125}
126
127/// List filesystem changes in a container.
128///
129/// # Examples
130///
131/// ```no_run
132/// use ironflow_ops_docker::containers::ContainerChanges;
133/// use ironflow_ops_docker::DockerClient;
134/// use ironflow_core::operation::Operation;
135///
136/// let client = DockerClient::connect_local().unwrap();
137/// let op = ContainerChanges::new(&client, "my-container");
138/// assert_eq!(op.kind(), "docker");
139/// ```
140pub struct ContainerChanges {
141    docker: Docker,
142    container: String,
143}
144
145impl ContainerChanges {
146    /// Create a new container-changes operation.
147    pub fn new(client: impl Into<DockerRef>, container: impl Into<String>) -> Self {
148        Self {
149            docker: client.into().0,
150            container: container.into(),
151        }
152    }
153
154    /// Execute and return a typed result.
155    ///
156    /// # Errors
157    ///
158    /// Returns [`OperationError::External`] if the container does not exist.
159    pub async fn run(
160        &self,
161        _ctx: &OperationContext,
162    ) -> Result<ContainerChangesOutput, OperationError> {
163        let response = self
164            .docker
165            .container_changes(&self.container)
166            .await
167            .map_err(docker_error)?;
168        let changes = response
169            .unwrap_or_default()
170            .into_iter()
171            .map(|c| ContainerChange {
172                path: c.path,
173                kind: c.kind as i32,
174            })
175            .collect();
176        Ok(ContainerChangesOutput { changes })
177    }
178}
179
180#[async_trait]
181impl Operation for ContainerChanges {
182    fn kind(&self) -> &str {
183        "docker"
184    }
185
186    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
187        to_value(&self.run(ctx).await?)
188    }
189
190    fn input(&self) -> Option<Value> {
191        Some(serde_json::json!({
192            "operation": "container_changes",
193            "container": self.container,
194        }))
195    }
196}
197
198impl TypedOperation for ContainerChanges {
199    type Output = ContainerChangesOutput;
200}
201
202// ---------------------------------------------------------------------------
203// ContainerPrune
204// ---------------------------------------------------------------------------
205
206/// Output of pruning stopped containers.
207#[derive(Debug, Clone, Serialize, Deserialize)]
208pub struct ContainerPruneOutput {
209    /// IDs of removed containers.
210    pub containers_deleted: Vec<String>,
211    /// Disk space reclaimed in bytes.
212    pub space_reclaimed: u64,
213}
214
215/// Remove all stopped containers.
216///
217/// # Examples
218///
219/// ```no_run
220/// use ironflow_ops_docker::containers::ContainerPrune;
221/// use ironflow_ops_docker::DockerClient;
222/// use ironflow_core::operation::Operation;
223///
224/// let client = DockerClient::connect_local().unwrap();
225/// let op = ContainerPrune::new(&client);
226/// assert_eq!(op.kind(), "docker");
227/// ```
228pub struct ContainerPrune {
229    docker: Docker,
230    filters: HashMap<String, Vec<String>>,
231}
232
233impl ContainerPrune {
234    /// Create a new container-prune operation.
235    pub fn new(client: impl Into<DockerRef>) -> Self {
236        Self {
237            docker: client.into().0,
238            filters: HashMap::new(),
239        }
240    }
241
242    /// Add a filter.
243    pub fn filter(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
244        self.filters
245            .entry(key.into())
246            .or_default()
247            .push(value.into());
248        self
249    }
250
251    /// Execute and return a typed result.
252    ///
253    /// # Errors
254    ///
255    /// Returns [`OperationError::External`] if the Docker daemon is
256    /// unreachable.
257    pub async fn run(
258        &self,
259        _ctx: &OperationContext,
260    ) -> Result<ContainerPruneOutput, OperationError> {
261        let options = PruneContainersOptions {
262            filters: Some(self.filters.clone()),
263        };
264        let response = self
265            .docker
266            .prune_containers(Some(options))
267            .await
268            .map_err(docker_error)?;
269        Ok(ContainerPruneOutput {
270            containers_deleted: response.containers_deleted.unwrap_or_default(),
271            space_reclaimed: response.space_reclaimed.unwrap_or(0) as u64,
272        })
273    }
274}
275
276#[async_trait]
277impl Operation for ContainerPrune {
278    fn kind(&self) -> &str {
279        "docker"
280    }
281
282    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
283        to_value(&self.run(ctx).await?)
284    }
285
286    fn input(&self) -> Option<Value> {
287        Some(serde_json::json!({
288            "operation": "container_prune",
289        }))
290    }
291}
292
293impl TypedOperation for ContainerPrune {
294    type Output = ContainerPruneOutput;
295}