1use std::future::Future;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use bollard::container::LogOutput;
6use bollard::exec::{StartExecOptions, StartExecResults};
7use bollard::models::ExecConfig;
8use futures_util::StreamExt;
9use soma_fleet::{HostRecord, TopologyRevision};
10use soma_ops::{MutationSendState, Timestamp};
11use tokio_util::sync::CancellationToken;
12
13use crate::{
14 BollardReadClient, ContainerExecMutator, ContainerExecReceipt, ContainerExecRequest,
15 InfraError, MutationFailure, MutationResult,
16};
17
18#[async_trait]
19impl ContainerExecMutator for BollardReadClient {
20 async fn exec_container(
21 &self,
22 host: &HostRecord,
23 request: &ContainerExecRequest,
24 cancellation: &CancellationToken,
25 ) -> MutationResult<ContainerExecReceipt> {
26 self.validate_host(host)
27 .map_err(|error| MutationFailure::new(MutationSendState::NotSent, error))?;
28 ensure_before_start(request.deadline(), cancellation)?;
29 let created = await_pre_start(
30 request.deadline(),
31 cancellation,
32 self.docker().create_exec(
33 request.container(),
34 ExecConfig {
35 cmd: Some(request.command().to_vec()),
36 user: request.user().map(str::to_owned),
37 working_dir: request
38 .working_dir()
39 .map(|path| path.to_string_lossy().into_owned()),
40 attach_stdout: Some(true),
41 attach_stderr: Some(true),
42 tty: Some(false),
43 ..Default::default()
44 },
45 ),
46 )
47 .await?;
48
49 let started = await_post_start(
50 request.deadline(),
51 cancellation,
52 self.docker().start_exec(
53 &created.id,
54 Some(StartExecOptions {
55 detach: false,
56 tty: false,
57 ..Default::default()
58 }),
59 ),
60 )
61 .await?;
62 let (mut stdout, mut stderr) = (Vec::new(), Vec::new());
63 let (mut stdout_truncated, mut stderr_truncated) = (false, false);
64 match started {
65 StartExecResults::Attached { mut output, .. } => loop {
66 let timeout = remaining(request.deadline(), MutationSendState::Unknown)?;
67 let next = tokio::select! {
68 () = cancellation.cancelled() => return Err(MutationFailure::new(
69 MutationSendState::Unknown,
70 soma_fleet::FleetError::Cancelled.into(),
71 )),
72 result = tokio::time::timeout(timeout, output.next()) => match result {
73 Err(_) => return Err(MutationFailure::new(
74 MutationSendState::Unknown,
75 soma_fleet::FleetError::DeadlineExceeded.into(),
76 )),
77 Ok(value) => value,
78 }
79 };
80 let Some(frame) = next else { break };
81 match frame.map_err(|error| {
82 MutationFailure::new(
83 MutationSendState::Unknown,
84 InfraError::Docker(error.to_string()),
85 )
86 })? {
87 LogOutput::StdOut { message } | LogOutput::Console { message } => {
88 stdout_truncated |=
89 append_bounded(&mut stdout, &message, request.max_stdout_bytes());
90 }
91 LogOutput::StdErr { message } => {
92 stderr_truncated |=
93 append_bounded(&mut stderr, &message, request.max_stderr_bytes());
94 }
95 _ => {}
96 }
97 },
98 StartExecResults::Detached => {
99 return Err(MutationFailure::new(
100 MutationSendState::Unknown,
101 InfraError::Docker(
102 "container exec unexpectedly detached; completion is unknown".into(),
103 ),
104 ));
105 }
106 }
107 let inspected = await_post_start(
108 request.deadline(),
109 cancellation,
110 self.docker().inspect_exec(&created.id),
111 )
112 .await?;
113 let stdout_text = String::from_utf8_lossy(&stdout);
114 let stderr_text = String::from_utf8_lossy(&stderr);
115 Ok(ContainerExecReceipt {
116 host: host.id().clone(),
117 topology_revision: TopologyRevision::clone(host.revision()),
118 container: request.container().to_owned(),
119 command: request.command().to_vec(),
120 user: request.user().map(str::to_owned),
121 working_dir: request.working_dir().map(ToOwned::to_owned),
122 stdout: stdout_text.into_owned(),
123 stderr: stderr_text.into_owned(),
124 exit_code: inspected.exit_code,
125 truncated: stdout_truncated || stderr_truncated,
126 encoding_lossy: std::str::from_utf8(&stdout).is_err()
127 || std::str::from_utf8(&stderr).is_err(),
128 send_state: MutationSendState::Sent,
129 })
130 }
131}
132
133fn ensure_before_start(
134 deadline: Timestamp,
135 cancellation: &CancellationToken,
136) -> MutationResult<()> {
137 if cancellation.is_cancelled() {
138 return Err(MutationFailure::new(
139 MutationSendState::NotSent,
140 soma_fleet::FleetError::Cancelled.into(),
141 ));
142 }
143 remaining(deadline, MutationSendState::NotSent).map(|_| ())
144}
145
146async fn await_pre_start<T, F>(
147 deadline: Timestamp,
148 cancellation: &CancellationToken,
149 future: F,
150) -> MutationResult<T>
151where
152 F: Future<Output = Result<T, bollard::errors::Error>>,
153{
154 let timeout = remaining(deadline, MutationSendState::NotSent)?;
155 tokio::select! {
156 () = cancellation.cancelled() => Err(MutationFailure::new(
157 MutationSendState::NotSent,
158 soma_fleet::FleetError::Cancelled.into(),
159 )),
160 result = tokio::time::timeout(timeout, future) => match result {
161 Err(_) => Err(MutationFailure::new(
162 MutationSendState::NotSent,
163 soma_fleet::FleetError::DeadlineExceeded.into(),
164 )),
165 Ok(Err(error)) => Err(MutationFailure::new(
166 MutationSendState::NotSent,
167 InfraError::Docker(error.to_string()),
168 )),
169 Ok(Ok(value)) => Ok(value),
170 }
171 }
172}
173
174async fn await_post_start<T, F>(
175 deadline: Timestamp,
176 cancellation: &CancellationToken,
177 future: F,
178) -> MutationResult<T>
179where
180 F: Future<Output = Result<T, bollard::errors::Error>>,
181{
182 let timeout = remaining(deadline, MutationSendState::Unknown)?;
183 tokio::select! {
184 () = cancellation.cancelled() => Err(MutationFailure::new(
185 MutationSendState::Unknown,
186 soma_fleet::FleetError::Cancelled.into(),
187 )),
188 result = tokio::time::timeout(timeout, future) => match result {
189 Err(_) => Err(MutationFailure::new(
190 MutationSendState::Unknown,
191 soma_fleet::FleetError::DeadlineExceeded.into(),
192 )),
193 Ok(Err(error)) => Err(MutationFailure::new(
194 MutationSendState::Unknown,
195 InfraError::Docker(error.to_string()),
196 )),
197 Ok(Ok(value)) => Ok(value),
198 }
199 }
200}
201
202fn remaining(deadline: Timestamp, send_state: MutationSendState) -> MutationResult<Duration> {
203 let millis = deadline
204 .unix_millis()
205 .saturating_sub(Timestamp::now().unix_millis());
206 if millis <= 0 {
207 Err(MutationFailure::new(
208 send_state,
209 soma_fleet::FleetError::DeadlineExceeded.into(),
210 ))
211 } else {
212 Ok(Duration::from_millis(millis as u64))
213 }
214}
215
216fn append_bounded(destination: &mut Vec<u8>, bytes: &[u8], limit: usize) -> bool {
217 let remaining = limit.saturating_sub(destination.len());
218 let retained = remaining.min(bytes.len());
219 destination.extend_from_slice(&bytes[..retained]);
220 retained < bytes.len()
221}
222
223#[cfg(test)]
224#[path = "bollard_exec_tests.rs"]
225mod tests;