Skip to main content

microsandbox_control_client/
json_client.rs

1//! Legacy unary exchanges: one real JSON operation per independently owned stream.
2
3use 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
21//--------------------------------------------------------------------------------------------------
22// Constants
23//--------------------------------------------------------------------------------------------------
24
25/// Bound applies to new-client discovery replies, not legacy request lines.
26pub const MAX_DISCOVERY_RESPONSE_SIZE: usize = 64 * 1024;
27
28//--------------------------------------------------------------------------------------------------
29// Types
30//--------------------------------------------------------------------------------------------------
31
32/// Explicit JSON adapter. Construction is inert; every call opens one stream.
33#[derive(Clone)]
34pub struct JsonControlClient {
35    pub(crate) dialer: Dialer,
36    pub(crate) options: ConnectOptions,
37    closed: CancellationToken,
38    rediscover: Option<ControlMode>,
39}
40
41//--------------------------------------------------------------------------------------------------
42// Methods
43//--------------------------------------------------------------------------------------------------
44
45impl JsonControlClient {
46    /// Select legacy JSON explicitly, without probing or opening a connection.
47    pub fn new(path: impl AsRef<Path>) -> Self {
48        Self::from_connector(Arc::new(LocalConnector::new(path)))
49    }
50
51    /// Store a repeatable dialer without I/O or hidden negotiation.
52    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    /// Configure local deadlines without performing I/O.
61    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    /// Inspect shared closure, including unexpected exchange failure.
88    pub fn is_closed(&self) -> bool {
89        self.closed.is_cancelled()
90    }
91
92    /// Wait for shared closure. Successful per-operation JSON EOF is not a
93    /// session closure; explicit close or an unexpected exchange failure is.
94    pub async fn closed(&self) {
95        self.closed.cancelled().await;
96    }
97
98    /// Close every clone and wake each outstanding exchange with its own
99    /// admission certainty. No implicit reconnect follows an explicit close.
100    pub async fn close(&self) {
101        self.closed.cancel();
102    }
103
104    /// Translate a known native request, returning its actual JSON reply.
105    pub async fn request(
106        &self,
107        message: impl IntoControlMessage,
108    ) -> ControlClientResult<JsonReply> {
109        self.request_with(message, |options| options).await
110    }
111
112    /// Configure one total attempt deadline, including dial and any discovery.
113    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    /// Normalize a checked operation without manufacturing a CBOR response.
124    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    /// Configure one checked unary attempt.
132    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    /// Execute a generation-aware checked request using its historical JSON representation.
147    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            // A path-only caller cannot reuse evidence about an old process.
183            // Discover again before this new exchange; a format change closes
184            // this handle instead of silently retargeting its prepared request.
185            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                    // Mark uncertainty before the first write poll: even a
249                    // partial failed write may have reached the peer.
250                    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
293//--------------------------------------------------------------------------------------------------
294// Functions
295//--------------------------------------------------------------------------------------------------
296
297pub(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}