Skip to main content

soma_infra/
bollard_exec.rs

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;