ironflow_ops_docker/containers/
cleanup.rs1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct ContainerStatsOutput {
24 pub data: Value,
26}
27
28pub struct ContainerStats {
42 docker: Docker,
43 container: String,
44}
45
46impl ContainerStats {
47 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 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#[derive(Debug, Clone, Serialize, Deserialize)]
113pub struct ContainerChange {
114 pub path: String,
116 pub kind: i32,
118}
119
120#[derive(Debug, Clone, Serialize, Deserialize)]
122pub struct ContainerChangesOutput {
123 pub changes: Vec<ContainerChange>,
125}
126
127pub struct ContainerChanges {
141 docker: Docker,
142 container: String,
143}
144
145impl ContainerChanges {
146 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 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#[derive(Debug, Clone, Serialize, Deserialize)]
208pub struct ContainerPruneOutput {
209 pub containers_deleted: Vec<String>,
211 pub space_reclaimed: u64,
213}
214
215pub struct ContainerPrune {
229 docker: Docker,
230 filters: HashMap<String, Vec<String>>,
231}
232
233impl ContainerPrune {
234 pub fn new(client: impl Into<DockerRef>) -> Self {
236 Self {
237 docker: client.into().0,
238 filters: HashMap::new(),
239 }
240 }
241
242 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 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}