ironflow_ops_docker/containers/
exec.rs1use async_trait::async_trait;
4use bollard::Docker;
5use bollard::exec::{CreateExecOptions, StartExecResults};
6use bollard::query_parameters::{LogsOptions, WaitContainerOptions};
7use ironflow_core::error::OperationError;
8use ironflow_core::operation::{Operation, OperationContext, TypedOperation};
9use serde::{Deserialize, Serialize};
10use serde_json::Value;
11use tokio_stream::StreamExt;
12
13use crate::containers::DockerRef;
14use crate::helpers::{docker_error, to_value};
15
16#[derive(Debug, Clone, Serialize, Deserialize)]
22pub struct ContainerLogsOutput {
23 pub lines: Vec<String>,
25}
26
27pub struct ContainerLogs {
41 docker: Docker,
42 container: String,
43 stdout: bool,
44 stderr: bool,
45 tail: Option<String>,
46}
47
48impl ContainerLogs {
49 pub fn new(client: impl Into<DockerRef>, container: impl Into<String>) -> Self {
51 Self {
52 docker: client.into().0,
53 container: container.into(),
54 stdout: true,
55 stderr: true,
56 tail: None,
57 }
58 }
59
60 pub fn stdout_only(mut self) -> Self {
62 self.stderr = false;
63 self
64 }
65
66 pub fn stderr_only(mut self) -> Self {
68 self.stdout = false;
69 self
70 }
71
72 pub fn tail(mut self, n: u64) -> Self {
74 self.tail = Some(n.to_string());
75 self
76 }
77
78 pub async fn run(
84 &self,
85 _ctx: &OperationContext,
86 ) -> Result<ContainerLogsOutput, OperationError> {
87 let options = LogsOptions {
88 stdout: self.stdout,
89 stderr: self.stderr,
90 tail: self.tail.clone().unwrap_or_else(|| "all".to_string()),
91 ..Default::default()
92 };
93 let mut stream = self.docker.logs(&self.container, Some(options));
94 let mut lines = Vec::new();
95 while let Some(result) = stream.next().await {
96 let chunk = result.map_err(docker_error)?;
97 lines.push(chunk.to_string());
98 }
99 Ok(ContainerLogsOutput { lines })
100 }
101}
102
103#[async_trait]
104impl Operation for ContainerLogs {
105 fn kind(&self) -> &str {
106 "docker"
107 }
108
109 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
110 to_value(&self.run(ctx).await?)
111 }
112
113 fn input(&self) -> Option<Value> {
114 Some(serde_json::json!({
115 "operation": "container_logs",
116 "container": self.container,
117 }))
118 }
119}
120
121impl TypedOperation for ContainerLogs {
122 type Output = ContainerLogsOutput;
123}
124
125#[derive(Debug, Clone, Serialize, Deserialize)]
131pub struct ContainerExecOutput {
132 pub output: Vec<String>,
134}
135
136pub struct ContainerExec {
150 docker: Docker,
151 container: String,
152 cmd: Vec<String>,
153}
154
155impl ContainerExec {
156 pub fn new(
158 client: impl Into<DockerRef>,
159 container: impl Into<String>,
160 cmd: Vec<impl Into<String>>,
161 ) -> Self {
162 Self {
163 docker: client.into().0,
164 container: container.into(),
165 cmd: cmd.into_iter().map(Into::into).collect(),
166 }
167 }
168
169 pub async fn run(
176 &self,
177 _ctx: &OperationContext,
178 ) -> Result<ContainerExecOutput, OperationError> {
179 let exec_options = CreateExecOptions {
180 cmd: Some(self.cmd.clone()),
181 attach_stdout: Some(true),
182 attach_stderr: Some(true),
183 ..Default::default()
184 };
185 let exec = self
186 .docker
187 .create_exec(&self.container, exec_options)
188 .await
189 .map_err(docker_error)?;
190 let start_result = self
191 .docker
192 .start_exec(&exec.id, None)
193 .await
194 .map_err(docker_error)?;
195 let mut output = Vec::new();
196 if let StartExecResults::Attached {
197 output: mut stream, ..
198 } = start_result
199 {
200 while let Some(result) = stream.next().await {
201 let chunk = result.map_err(docker_error)?;
202 output.push(chunk.to_string());
203 }
204 }
205 Ok(ContainerExecOutput { output })
206 }
207}
208
209#[async_trait]
210impl Operation for ContainerExec {
211 fn kind(&self) -> &str {
212 "docker"
213 }
214
215 async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
216 to_value(&self.run(ctx).await?)
217 }
218
219 fn input(&self) -> Option<Value> {
220 Some(serde_json::json!({
221 "operation": "container_exec",
222 "container": self.container,
223 "cmd": self.cmd,
224 }))
225 }
226}
227
228impl TypedOperation for ContainerExec {
229 type Output = ContainerExecOutput;
230}
231
232#[derive(Debug, Clone, Serialize, Deserialize)]
238pub struct ContainerWaitOutput {
239 pub status_code: i64,
241}
242
243pub struct ContainerWait {
257 docker: Docker,
258 container: String,
259}
260
261impl ContainerWait {
262 pub fn new(client: impl Into<DockerRef>, container: impl Into<String>) -> Self {
264 Self {
265 docker: client.into().0,
266 container: container.into(),
267 }
268 }
269
270 pub async fn run(
276 &self,
277 _ctx: &OperationContext,
278 ) -> Result<ContainerWaitOutput, OperationError> {
279 let options = WaitContainerOptions {
280 condition: "not-running".to_string(),
281 };
282 let mut stream = self.docker.wait_container(&self.container, Some(options));
283 let mut status_code = 0i64;
284 while let Some(result) = stream.next().await {
285 let response = result.map_err(docker_error)?;
286 status_code = response.status_code;
287 }
288 Ok(ContainerWaitOutput { status_code })
289 }
290}
291
292#[async_trait]
293impl Operation for ContainerWait {
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": "container_wait",
305 "container": self.container,
306 }))
307 }
308}
309
310impl TypedOperation for ContainerWait {
311 type Output = ContainerWaitOutput;
312}