soma_infra/
bollard_mutation.rs1use std::future::Future;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use bollard::query_parameters::{
6 RestartContainerOptions, StartContainerOptions, StopContainerOptions,
7};
8use soma_fleet::{HostRecord, TopologyRevision};
9use soma_ops::{MutationSendState, Timestamp};
10use tokio_util::sync::CancellationToken;
11
12use crate::{
13 BollardReadClient, ContainerLifecycleAction, ContainerLifecycleMutator,
14 ContainerLifecycleRequest, ContainerMutationReceipt, InfraError, MutationFailure,
15 MutationResult,
16};
17
18#[async_trait]
19impl ContainerLifecycleMutator for BollardReadClient {
20 async fn mutate_container(
21 &self,
22 host: &HostRecord,
23 request: &ContainerLifecycleRequest,
24 cancellation: &CancellationToken,
25 ) -> MutationResult<ContainerMutationReceipt> {
26 self.validate_host(host)
27 .map_err(|error| MutationFailure::new(MutationSendState::NotSent, error))?;
28 ensure_not_expired(request.deadline(), cancellation)?;
29 let result = match request.action() {
30 ContainerLifecycleAction::Start => {
31 await_send(
32 request.deadline(),
33 cancellation,
34 self.docker()
35 .start_container(request.container(), None::<StartContainerOptions>),
36 )
37 .await
38 }
39 ContainerLifecycleAction::Stop => {
40 await_send(
41 request.deadline(),
42 cancellation,
43 self.docker()
44 .stop_container(request.container(), None::<StopContainerOptions>),
45 )
46 .await
47 }
48 ContainerLifecycleAction::Restart => {
49 await_send(
50 request.deadline(),
51 cancellation,
52 self.docker()
53 .restart_container(request.container(), None::<RestartContainerOptions>),
54 )
55 .await
56 }
57 ContainerLifecycleAction::Pause => {
58 await_send(
59 request.deadline(),
60 cancellation,
61 self.docker().pause_container(request.container()),
62 )
63 .await
64 }
65 ContainerLifecycleAction::Resume => {
66 await_send(
67 request.deadline(),
68 cancellation,
69 self.docker().unpause_container(request.container()),
70 )
71 .await
72 }
73 };
74 result?;
75 Ok(ContainerMutationReceipt {
76 host: host.id().clone(),
77 topology_revision: TopologyRevision::clone(host.revision()),
78 container: request.container().to_owned(),
79 action: request.action(),
80 send_state: MutationSendState::Sent,
81 })
82 }
83}
84
85fn ensure_not_expired(deadline: Timestamp, cancellation: &CancellationToken) -> MutationResult<()> {
86 if cancellation.is_cancelled() {
87 return Err(MutationFailure::new(
88 MutationSendState::NotSent,
89 soma_fleet::FleetError::Cancelled.into(),
90 ));
91 }
92 if Timestamp::now() >= deadline {
93 return Err(MutationFailure::new(
94 MutationSendState::NotSent,
95 soma_fleet::FleetError::DeadlineExceeded.into(),
96 ));
97 }
98 Ok(())
99}
100
101async fn await_send<F>(
102 deadline: Timestamp,
103 cancellation: &CancellationToken,
104 future: F,
105) -> MutationResult<()>
106where
107 F: Future<Output = Result<(), bollard::errors::Error>>,
108{
109 let now = Timestamp::now().unix_millis();
110 let remaining = deadline.unix_millis().saturating_sub(now);
111 if remaining <= 0 {
112 return Err(MutationFailure::new(
113 MutationSendState::NotSent,
114 soma_fleet::FleetError::DeadlineExceeded.into(),
115 ));
116 }
117 tokio::select! {
118 () = cancellation.cancelled() => Err(MutationFailure::new(
119 MutationSendState::Unknown,
120 soma_fleet::FleetError::Cancelled.into(),
121 )),
122 () = tokio::time::sleep(Duration::from_millis(remaining as u64)) => Err(
123 MutationFailure::new(
124 MutationSendState::Unknown,
125 soma_fleet::FleetError::DeadlineExceeded.into(),
126 )
127 ),
128 result = future => result.map_err(|error| MutationFailure::new(
129 MutationSendState::Unknown,
130 InfraError::Docker(error.to_string()),
131 )),
132 }
133}
134
135#[cfg(test)]
136#[path = "bollard_mutation_tests.rs"]
137mod tests;