microsandbox_control_client/
json_client.rs1use std::path::Path;
4use std::sync::Arc;
5use std::time::Duration;
6
7use microsandbox_protocol::control::{ControlRequest, DEFAULT_REQUEST_TIMEOUT};
8use microsandbox_protocol_client::{
9 ClientError, ConnectOptions, Connector, Delivery, ErrorKind, LocalConnector, RequestOptions,
10};
11use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
12use tokio::time::{Instant, timeout_at};
13use tokio_util::sync::CancellationToken;
14use zeroize::Zeroizing;
15
16use crate::{
17 CheckedControlRequest, CompatibleControlRequest, ControlClientError, ControlClientResult,
18 ControlMode, IntoControlMessage, JsonReply, dialer::Dialer,
19};
20
21pub const MAX_DISCOVERY_RESPONSE_SIZE: usize = 64 * 1024;
27
28#[derive(Clone)]
34pub struct JsonControlClient {
35 pub(crate) dialer: Dialer,
36 pub(crate) options: ConnectOptions,
37 closed: CancellationToken,
38 rediscover: Option<ControlMode>,
39}
40
41impl JsonControlClient {
46 pub fn new(path: impl AsRef<Path>) -> Self {
48 Self::from_connector(Arc::new(LocalConnector::new(path)))
49 }
50
51 pub fn from_connector(connector: Arc<dyn Connector>) -> Self {
53 Self::configured(
54 Dialer::Unverified(connector),
55 ConnectOptions::default(),
56 None,
57 )
58 }
59
60 pub fn from_connector_with(
62 connector: Arc<dyn Connector>,
63 configure: impl FnOnce(ConnectOptions) -> ConnectOptions,
64 ) -> ControlClientResult<Self> {
65 let options = configure(ConnectOptions::default());
66 options.limits.validate()?;
67 Ok(Self::configured(
68 Dialer::Unverified(connector),
69 options,
70 None,
71 ))
72 }
73
74 pub(crate) fn configured(
75 dialer: Dialer,
76 options: ConnectOptions,
77 rediscover: Option<ControlMode>,
78 ) -> Self {
79 Self {
80 dialer,
81 options,
82 closed: CancellationToken::new(),
83 rediscover,
84 }
85 }
86
87 pub fn is_closed(&self) -> bool {
89 self.closed.is_cancelled()
90 }
91
92 pub async fn closed(&self) {
95 self.closed.cancelled().await;
96 }
97
98 pub async fn close(&self) {
101 self.closed.cancel();
102 }
103
104 pub async fn request(
106 &self,
107 message: impl IntoControlMessage,
108 ) -> ControlClientResult<JsonReply> {
109 self.request_with(message, |options| options).await
110 }
111
112 pub async fn request_with(
114 &self,
115 message: impl IntoControlMessage,
116 configure: impl FnOnce(RequestOptions) -> RequestOptions,
117 ) -> ControlClientResult<JsonReply> {
118 let request = message.into_json()?;
119 self.operation(request, configure(RequestOptions::default()))
120 .await
121 }
122
123 pub async fn request_typed<R: CheckedControlRequest>(
125 &self,
126 request: &R,
127 ) -> ControlClientResult<R::Response> {
128 self.request_typed_with(request, |options| options).await
129 }
130
131 pub async fn request_typed_with<R: CheckedControlRequest>(
133 &self,
134 request: &R,
135 configure: impl FnOnce(RequestOptions) -> RequestOptions,
136 ) -> ControlClientResult<R::Response> {
137 let response = self
138 .operation(
139 request.json_request()?,
140 configure(RequestOptions::default()),
141 )
142 .await?;
143 request.decode_json(response)
144 }
145
146 pub async fn request_compatible<R: CompatibleControlRequest>(
148 &self,
149 request: &R,
150 options: RequestOptions,
151 ) -> ControlClientResult<R::Response> {
152 let response = self
153 .operation_bytes(request.compatibility_json_bytes()?, options)
154 .await?;
155 request.decode_compatibility_json(response)
156 }
157
158 pub(crate) async fn operation(
159 &self,
160 request: ControlRequest,
161 options: RequestOptions,
162 ) -> ControlClientResult<JsonReply> {
163 let bytes = Zeroizing::new(
164 serde_json::to_vec(&request).map_err(|_| ClientError::new(ErrorKind::Encode))?,
165 );
166 self.operation_bytes(bytes, options).await
167 }
168
169 pub(crate) async fn operation_bytes(
170 &self,
171 request: Zeroizing<Vec<u8>>,
172 options: RequestOptions,
173 ) -> ControlClientResult<JsonReply> {
174 let until = deadline(
175 options
176 .request_timeout
177 .or(self.options.limits.request_timeout)
178 .unwrap_or(DEFAULT_REQUEST_TIMEOUT),
179 )?;
180 let setup_until = deadline(self.options.setup_timeout)?.min(until);
181 if let Some(expected) = self.rediscover.filter(|_| !self.dialer.verified()) {
182 let mode = self
186 .discover(setup_until)
187 .await
188 .map_err(crate::connection::not_sent);
189 let mode = match mode {
190 Ok((mode, _)) => mode,
191 Err(error) => {
192 self.close().await;
193 return Err(error);
194 }
195 };
196 if mode != expected {
197 self.close().await;
198 return Err(ControlClientError::RuntimeChanged);
199 }
200 }
201 self.exchange_bytes(request, until, setup_until, None).await
202 }
203
204 pub(crate) async fn discover(
205 &self,
206 until: Instant,
207 ) -> ControlClientResult<(ControlMode, crate::RuntimeCapabilities)> {
208 let request = Zeroizing::new(
209 serde_json::to_vec(&ControlRequest::Capabilities)
210 .map_err(|_| ClientError::new(ErrorKind::Encode))?,
211 );
212 let reply = self
213 .exchange_bytes(request, until, until, Some(MAX_DISCOVERY_RESPONSE_SIZE))
214 .await?;
215 let mode = reply.discovery_mode()?;
216 let capabilities = crate::json_reply::capabilities(
217 reply
218 .value()
219 .get("capabilities")
220 .ok_or_else(|| ClientError::new(ErrorKind::InvalidData))?,
221 )
222 .ok_or_else(|| ClientError::new(ErrorKind::InvalidData))?;
223 Ok((mode, capabilities))
224 }
225
226 async fn exchange_bytes(
227 &self,
228 mut line: Zeroizing<Vec<u8>>,
229 until: Instant,
230 setup_until: Instant,
231 max_reply: Option<usize>,
232 ) -> ControlClientResult<JsonReply> {
233 line.push(b'\n');
234 let mut admitted = false;
235 let result = timeout_at(until, async {
236 tokio::select! {
237 biased;
238 _ = self.closed.cancelled() => Err(ClientError::new(ErrorKind::Closed).into()),
239 result = async {
240 check_deadline(setup_until)?;
241 let mut transport = timeout_at(setup_until, self.dialer.connect(setup_until)).await
242 .map_err(|_| ClientError::new(ErrorKind::Timeout))??;
243 check_deadline(setup_until)?;
244 timeout_at(setup_until, self.dialer.verify(setup_until)).await
245 .map_err(|_| ClientError::new(ErrorKind::Timeout))??;
246 check_deadline(until)?;
247 if self.closed.is_cancelled() { return Err(ClientError::new(ErrorKind::Closed).into()); }
248 admitted = true;
251 transport.write_all(&line).await.map_err(ClientError::from)?;
252 transport.flush().await.map_err(ClientError::from)?;
253 let mut reader = BufReader::new(transport);
254 let mut reply = Vec::new();
255 loop {
256 let buffer = reader.fill_buf().await.map_err(ClientError::from)?;
257 if buffer.is_empty() {
258 if reply.is_empty() { return Err(ClientError::new(ErrorKind::PeerClosed).into()); }
259 break;
260 }
261 let end = buffer.iter().position(|byte| *byte == b'\n').map(|index| index + 1);
262 let count = end.unwrap_or(buffer.len());
263 if max_reply.is_some_and(|limit| reply.len().saturating_add(count) > limit) {
264 return Err(ClientError::new(ErrorKind::InvalidData).into());
265 }
266 reply.extend_from_slice(&buffer[..count]);
267 reader.consume(count);
268 if end.is_some() { break; }
269 }
270 JsonReply::parse(reply)
271 } => result,
272 }
273 }).await.unwrap_or_else(|_| Err(ClientError::new(ErrorKind::Timeout).into()));
274 match result {
275 Ok(reply) => Ok(reply),
276 Err(error) => {
277 self.closed.cancel();
278 Err(match error {
279 ControlClientError::Client(error) => error
280 .with_delivery(if admitted {
281 Delivery::Unknown
282 } else {
283 Delivery::NotSent
284 })
285 .into(),
286 error => error,
287 })
288 }
289 }
290 }
291}
292
293pub(crate) fn deadline(duration: Duration) -> ControlClientResult<Instant> {
298 Instant::now()
299 .checked_add(duration)
300 .ok_or_else(|| ClientError::new(ErrorKind::InvalidOptions).into())
301}
302
303pub(crate) fn check_deadline(until: Instant) -> ControlClientResult<()> {
304 if Instant::now() >= until {
305 Err(ClientError::new(ErrorKind::Timeout).into())
306 } else {
307 Ok(())
308 }
309}