Skip to main content

a3s_box_runtime/grpc/
exec.rs

1//! Command execution and streaming clients.
2
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5
6use a3s_box_core::error::{BoxError, Result};
7use a3s_box_core::ExecutionProcessSignal;
8use tokio::io::AsyncWriteExt;
9use tokio::sync::Mutex;
10
11#[cfg(unix)]
12type ExecStream = tokio::net::UnixStream;
13#[cfg(windows)]
14type ExecStream = tokio::net::windows::named_pipe::NamedPipeClient;
15
16const EXEC_CONTROL_CANCEL: &[u8] = b"cancel";
17const EXEC_CONTROL_STDIN_CLOSE: &[u8] = b"stdin-close";
18const EXEC_CONTROL_SIGNAL: &[u8] = b"signal:";
19/// Host→guest control: flush all buffered output and reply with a flush-ack.
20const EXEC_CONTROL_FLUSH: &[u8] = b"flush";
21/// Host→guest control: stream a tar archive of the guest-visible rootfs.
22const EXEC_CONTROL_ARCHIVE_ROOTFS: &[u8] = b"archive-rootfs-v1";
23const EXEC_CONTROL_ARCHIVE_ROOTFS_PAUSE: &[u8] = b"archive-rootfs-v1:pause";
24/// Guest→host marker after every archive data frame has been sent.
25const EXEC_ARCHIVE_ROOTFS_DONE: &[u8] = b"archive-rootfs-v1-done";
26/// Host→trusted-maintenance-PID1 request to unmount its rootfs disk and exit.
27const EXEC_CONTROL_SHUTDOWN_MAINTENANCE: &[u8] = b"shutdown-rootfs-maintenance-v1";
28const EXEC_SHUTDOWN_MAINTENANCE_ACK: &[u8] = b"shutdown-rootfs-maintenance-v1-ack";
29/// Guest→host marker (carried in a Control frame) acknowledging a flush. Kept
30/// distinct from an `ExecExit` JSON payload so `next_event` can tell them apart.
31/// Must match the guest's `EXEC_FLUSH_ACK` in `guest/init/src/exec_server.rs`.
32const EXEC_FLUSH_ACK: &[u8] = b"flush-ack";
33/// Guest→host acknowledgement that a `signal-main:<N>` graceful-stop control was
34/// received and the signal delivered. Must match the guest's
35/// `EXEC_SIGNAL_MAIN_ACK` in `guest/init/src/exec_server.rs`.
36const EXEC_SIGNAL_MAIN_ACK: &[u8] = b"signal-main-ack";
37/// Guest→host acknowledgement that a `spawn-main` deferred-main control was
38/// received and the container main spawned. Matches the guest's
39/// `EXEC_SPAWN_MAIN_ACK` in `guest/init/src/exec_server.rs`.
40const EXEC_SPAWN_MAIN_ACK: &[u8] = b"spawn-main-ack";
41/// Guest→host negative acknowledgement for `spawn-main`, followed by a UTF-8-ish
42/// diagnostic string from guest-init.
43const EXEC_SPAWN_MAIN_NACK: &[u8] = b"spawn-main-nack:";
44
45/// Host-side slack added to a one-shot exec's in-guest `timeout_ns` before the
46/// host gives up reading the reply. The in-guest timeout cannot fire if the
47/// guest is wedged, so the host needs its own ceiling.
48const EXEC_HOST_SLACK_SECS: u64 = 10;
49/// Host-side deadline for a `signal-main` ACK. Signal delivery + the ACK are
50/// fast; a wedged guest that never replies must not block the caller's
51/// force-kill fallback.
52const SIGNAL_MAIN_ACK_TIMEOUT_SECS: u64 = 10;
53
54type ExecFrameReader = a3s_transport::FrameReader<tokio::io::ReadHalf<ExecStream>>;
55type ExecFrameWriter = a3s_transport::FrameWriter<tokio::io::WriteHalf<ExecStream>>;
56
57async fn connect_exec_stream(path: &Path) -> std::io::Result<ExecStream> {
58    #[cfg(unix)]
59    {
60        tokio::net::UnixStream::connect(path).await
61    }
62
63    #[cfg(windows)]
64    {
65        const ERROR_FILE_NOT_FOUND: i32 = 2;
66        const ERROR_PIPE_BUSY: i32 = 231;
67        const CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
68        const RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(10);
69
70        let deadline = tokio::time::Instant::now() + CONNECT_TIMEOUT;
71        loop {
72            match tokio::net::windows::named_pipe::ClientOptions::new().open(path) {
73                Ok(stream) => return Ok(stream),
74                Err(error)
75                    if matches!(
76                        error.raw_os_error(),
77                        Some(ERROR_FILE_NOT_FOUND | ERROR_PIPE_BUSY)
78                    ) && tokio::time::Instant::now() < deadline =>
79                {
80                    tokio::time::sleep(RETRY_DELAY).await;
81                }
82                Err(error) => return Err(error),
83            }
84        }
85    }
86}
87
88/// Client for executing commands through the platform-local guest channel.
89///
90/// Uses the Frame wire protocol: sends a Data frame with JSON ExecRequest,
91/// receives a Data frame with JSON ExecOutput.
92#[derive(Debug)]
93pub struct ExecClient {
94    socket_path: PathBuf,
95}
96
97impl ExecClient {
98    pub(crate) fn for_socket(socket_path: &Path) -> Self {
99        Self {
100            socket_path: socket_path.to_path_buf(),
101        }
102    }
103
104    /// Connect to the exec server via a Unix socket or Windows named pipe.
105    ///
106    /// Verifies the socket is connectable.
107    pub async fn connect(socket_path: &Path) -> Result<Self> {
108        let client = Self::for_socket(socket_path);
109        let _stream = client.open_stream().await?;
110        Ok(client)
111    }
112
113    /// Get the socket path this client is connected to.
114    pub fn socket_path(&self) -> &Path {
115        &self.socket_path
116    }
117
118    pub(crate) async fn open_stream(&self) -> Result<ExecStream> {
119        connect_exec_stream(&self.socket_path).await.map_err(|e| {
120            BoxError::ExecError(format!(
121                "Exec connection failed to {}: {}",
122                self.socket_path.display(),
123                e,
124            ))
125        })
126    }
127
128    /// Execute a command in the guest.
129    ///
130    /// Sends a Data frame with JSON ExecRequest, reads a Data frame with JSON ExecOutput.
131    pub async fn exec_command(
132        &self,
133        request: &a3s_box_core::exec::ExecRequest,
134    ) -> Result<a3s_box_core::exec::ExecOutput> {
135        let stream = self.open_stream().await?;
136        self.exec_command_on_stream(stream, request).await
137    }
138
139    pub(crate) async fn exec_command_on_stream(
140        &self,
141        mut stream: ExecStream,
142        request: &a3s_box_core::exec::ExecRequest,
143    ) -> Result<a3s_box_core::exec::ExecOutput> {
144        let payload = serde_json::to_vec(request)
145            .map_err(|e| BoxError::ExecError(format!("Failed to serialize exec request: {}", e)))?;
146
147        // Send request as Data frame
148        let request_frame = a3s_transport::Frame::data(payload);
149        let encoded = request_frame.encode().map_err(|e| {
150            BoxError::ExecError(format!("Failed to encode exec request frame: {}", e))
151        })?;
152        stream
153            .write_all(&encoded)
154            .await
155            .map_err(|e| BoxError::ExecError(format!("Exec request write failed: {}", e)))?;
156
157        // Read response frame, bounded by a HOST-side deadline of the request's
158        // timeout plus slack. The request's timeout_ns is only enforced INSIDE
159        // the guest; a wedged guest (kernel hang, OOM thrash, frozen VM) can
160        // still complete the host connect handshake but never reply, which would
161        // block this read forever and stall every caller (health probes, the
162        // monitor poll loop, CLI exec).
163        let (r, _w) = tokio::io::split(stream);
164        let mut reader = a3s_transport::FrameReader::new(r);
165        let host_deadline = std::time::Duration::from_nanos(request.timeout_ns)
166            .saturating_add(std::time::Duration::from_secs(EXEC_HOST_SLACK_SECS));
167        let frame = tokio::time::timeout(host_deadline, reader.read_frame())
168            .await
169            .map_err(|_| {
170                BoxError::ExecError(format!(
171                    "Exec response timed out after {host_deadline:?} (guest may be wedged)"
172                ))
173            })?
174            .map_err(|e| BoxError::ExecError(format!("Exec response read failed: {}", e)))?
175            .ok_or_else(|| {
176                BoxError::ExecError("Exec server closed without response".to_string())
177            })?;
178
179        match frame.frame_type {
180            a3s_transport::FrameType::Data => {
181                let output: a3s_box_core::exec::ExecOutput = serde_json::from_slice(&frame.payload)
182                    .map_err(|e| {
183                        BoxError::ExecError(format!("Failed to parse exec response: {}", e))
184                    })?;
185                Ok(output)
186            }
187            a3s_transport::FrameType::Error => {
188                let msg = String::from_utf8_lossy(&frame.payload);
189                Err(BoxError::ExecError(format!("Exec server error: {}", msg)))
190            }
191            _ => Err(BoxError::ExecError(format!(
192                "Unexpected frame type: {:?}",
193                frame.frame_type
194            ))),
195        }
196    }
197
198    /// Execute a command in streaming mode.
199    ///
200    /// Sends a Data frame with JSON ExecRequest (streaming=true), then reads
201    /// multiple frames: ExecChunk frames for stdout/stderr data, and a final
202    /// ExecExit frame with the exit code.
203    ///
204    /// Returns a `StreamingExec` handle for reading events.
205    pub async fn exec_stream(
206        &self,
207        request: &a3s_box_core::exec::ExecRequest,
208    ) -> Result<StreamingExec> {
209        let stream = self.open_stream().await?;
210        self.exec_stream_on_stream(stream, request).await
211    }
212
213    pub(crate) async fn exec_stream_on_stream(
214        &self,
215        stream: ExecStream,
216        request: &a3s_box_core::exec::ExecRequest,
217    ) -> Result<StreamingExec> {
218        let mut req = request.clone();
219        req.streaming = true;
220
221        let payload = serde_json::to_vec(&req)
222            .map_err(|e| BoxError::ExecError(format!("Failed to serialize exec request: {}", e)))?;
223
224        let (r, w) = tokio::io::split(stream);
225        let mut writer = a3s_transport::FrameWriter::new(w);
226        writer
227            .write_data(&payload)
228            .await
229            .map_err(|e| BoxError::ExecError(format!("Exec request write failed: {}", e)))?;
230
231        let reader = a3s_transport::FrameReader::new(r);
232        let started = std::time::Instant::now();
233
234        Ok(StreamingExec {
235            reader,
236            writer: Arc::new(Mutex::new(writer)),
237            started,
238            stdout_bytes: 0,
239            stderr_bytes: 0,
240            done: false,
241        })
242    }
243
244    /// Stream a guest-created rootfs tar archive into `output`.
245    ///
246    /// The guest performs `stat` and tar-header creation, preserving Linux
247    /// uid/gid/mode even when the host virtio-fs backing directory exposes
248    /// different macOS metadata. Mounted subtrees are excluded by guest-init.
249    pub async fn archive_rootfs<W>(&self, output: &mut W, pause: bool) -> Result<u64>
250    where
251        W: tokio::io::AsyncWrite + Unpin,
252    {
253        let mut stream = connect_exec_stream(&self.socket_path)
254            .await
255            .map_err(|error| {
256                BoxError::ExecError(format!(
257                    "Rootfs archive connection failed to {}: {error}",
258                    self.socket_path.display()
259                ))
260            })?;
261
262        let control = if pause {
263            EXEC_CONTROL_ARCHIVE_ROOTFS_PAUSE
264        } else {
265            EXEC_CONTROL_ARCHIVE_ROOTFS
266        };
267        let request = a3s_transport::Frame::control(control.to_vec());
268        stream
269            .write_all(&request.encode().map_err(|error| {
270                BoxError::ExecError(format!("Rootfs archive request encode failed: {error}"))
271            })?)
272            .await
273            .map_err(|error| {
274                BoxError::ExecError(format!("Rootfs archive request write failed: {error}"))
275            })?;
276
277        let (reader, _writer) = tokio::io::split(stream);
278        let mut reader = a3s_transport::FrameReader::new(reader);
279        let mut written = 0u64;
280        loop {
281            let frame = reader
282                .read_frame()
283                .await
284                .map_err(|error| {
285                    BoxError::ExecError(format!("Rootfs archive read failed: {error}"))
286                })?
287                .ok_or_else(|| {
288                    BoxError::ExecError(
289                        "Rootfs archive stream closed before completion".to_string(),
290                    )
291                })?;
292
293            match frame.frame_type {
294                a3s_transport::FrameType::Data => {
295                    output.write_all(&frame.payload).await.map_err(|error| {
296                        BoxError::ExecError(format!("Rootfs archive output write failed: {error}"))
297                    })?;
298                    written = written.saturating_add(frame.payload.len() as u64);
299                }
300                a3s_transport::FrameType::Control if frame.payload == EXEC_ARCHIVE_ROOTFS_DONE => {
301                    output.flush().await.map_err(|error| {
302                        BoxError::ExecError(format!("Rootfs archive output flush failed: {error}"))
303                    })?;
304                    return Ok(written);
305                }
306                a3s_transport::FrameType::Error => {
307                    return Err(BoxError::ExecError(format!(
308                        "Guest rootfs archive failed: {}",
309                        String::from_utf8_lossy(&frame.payload)
310                    )));
311                }
312                other => {
313                    return Err(BoxError::ExecError(format!(
314                        "Unexpected rootfs archive frame: {other:?}"
315                    )));
316                }
317            }
318        }
319    }
320
321    /// Transfer a file to/from the guest.
322    ///
323    /// Sends a discriminated JSON file request and reads a JSON FileResponse.
324    pub async fn file_transfer(
325        &self,
326        request: &a3s_box_core::exec::FileRequest,
327    ) -> Result<a3s_box_core::exec::FileResponse> {
328        let stream = self.open_stream().await?;
329        self.file_transfer_on_stream(stream, request).await
330    }
331
332    pub(crate) async fn file_transfer_on_stream(
333        &self,
334        mut stream: ExecStream,
335        request: &a3s_box_core::exec::FileRequest,
336    ) -> Result<a3s_box_core::exec::FileResponse> {
337        let payload = serde_json::to_vec(&a3s_box_core::GuestSessionRequest::File(request.clone()))
338            .map_err(|e| BoxError::ExecError(format!("Failed to serialize file request: {}", e)))?;
339
340        let request_frame = a3s_transport::Frame::data(payload);
341        let encoded = request_frame.encode().map_err(|e| {
342            BoxError::ExecError(format!("Failed to encode file request frame: {}", e))
343        })?;
344        stream
345            .write_all(&encoded)
346            .await
347            .map_err(|e| BoxError::ExecError(format!("File request write failed: {}", e)))?;
348
349        let (r, _w) = tokio::io::split(stream);
350        let mut reader = a3s_transport::FrameReader::new(r);
351        let frame = reader
352            .read_frame()
353            .await
354            .map_err(|e| BoxError::ExecError(format!("File response read failed: {}", e)))?
355            .ok_or_else(|| {
356                BoxError::ExecError("Exec server closed without response".to_string())
357            })?;
358
359        match frame.frame_type {
360            a3s_transport::FrameType::Data => {
361                let response: a3s_box_core::exec::FileResponse =
362                    serde_json::from_slice(&frame.payload).map_err(|e| {
363                        BoxError::ExecError(format!("Failed to parse file response: {}", e))
364                    })?;
365                Ok(response)
366            }
367            a3s_transport::FrameType::Error => {
368                let msg = String::from_utf8_lossy(&frame.payload);
369                Err(BoxError::ExecError(format!("File transfer error: {}", msg)))
370            }
371            _ => Err(BoxError::ExecError(format!(
372                "Unexpected frame type: {:?}",
373                frame.frame_type
374            ))),
375        }
376    }
377
378    /// Perform a filesystem metadata or mutation operation inside the guest.
379    pub async fn filesystem(
380        &self,
381        request: &a3s_box_core::FilesystemRequest,
382    ) -> Result<a3s_box_core::FilesystemResponse> {
383        let stream = self.open_stream().await?;
384        self.filesystem_on_stream(stream, request).await
385    }
386
387    pub(crate) async fn filesystem_on_stream(
388        &self,
389        mut stream: ExecStream,
390        request: &a3s_box_core::FilesystemRequest,
391    ) -> Result<a3s_box_core::FilesystemResponse> {
392        let payload = serde_json::to_vec(&a3s_box_core::GuestSessionRequest::Filesystem(
393            request.clone(),
394        ))
395        .map_err(|error| {
396            BoxError::ExecError(format!("Failed to serialize filesystem request: {error}"))
397        })?;
398        let encoded = a3s_transport::Frame::data(payload)
399            .encode()
400            .map_err(|error| {
401                BoxError::ExecError(format!("Failed to encode filesystem request: {error}"))
402            })?;
403        stream.write_all(&encoded).await.map_err(|error| {
404            BoxError::ExecError(format!("Filesystem request write failed: {error}"))
405        })?;
406
407        let (read, _write) = tokio::io::split(stream);
408        let mut reader = a3s_transport::FrameReader::new(read);
409        let frame = reader
410            .read_frame()
411            .await
412            .map_err(|error| {
413                BoxError::ExecError(format!("Filesystem response read failed: {error}"))
414            })?
415            .ok_or_else(|| {
416                BoxError::ExecError("Exec server closed without filesystem response".to_string())
417            })?;
418        match frame.frame_type {
419            a3s_transport::FrameType::Data => {
420                serde_json::from_slice(&frame.payload).map_err(|error| {
421                    BoxError::ExecError(format!("Failed to parse filesystem response: {error}"))
422                })
423            }
424            a3s_transport::FrameType::Error => Err(BoxError::ExecError(format!(
425                "Filesystem operation failed: {}",
426                String::from_utf8_lossy(&frame.payload)
427            ))),
428            other => Err(BoxError::ExecError(format!(
429                "Unexpected filesystem response frame: {other:?}"
430            ))),
431        }
432    }
433
434    /// Send a Heartbeat frame and wait for a Heartbeat response.
435    ///
436    /// Returns `true` if the exec server responds, `false` otherwise.
437    pub async fn heartbeat(&self) -> Result<bool> {
438        let mut stream = match connect_exec_stream(&self.socket_path).await {
439            Ok(s) => s,
440            Err(_) => return Ok(false),
441        };
442
443        let frame = a3s_transport::Frame::heartbeat();
444        let encoded = match frame.encode() {
445            Ok(e) => e,
446            Err(_) => return Ok(false),
447        };
448
449        if stream.write_all(&encoded).await.is_err() {
450            return Ok(false);
451        }
452
453        let (r, _w) = tokio::io::split(stream);
454        let mut reader = a3s_transport::FrameReader::new(r);
455        match reader.read_frame().await {
456            Ok(Some(f)) if f.frame_type == a3s_transport::FrameType::Heartbeat => Ok(true),
457            _ => Ok(false),
458        }
459    }
460
461    /// Ask the guest to deliver `signal` (a signal number, e.g. 15 for SIGTERM)
462    /// to the main container process for graceful shutdown. The guest runs the
463    /// container's own stop handler; when it exits, guest init exits and the VM
464    /// stops cleanly. Returns `Ok(true)` if the guest acknowledged, `Ok(false)`
465    /// if it did not respond (caller should fall back to a hard stop).
466    pub async fn signal_main(&self, signal: i32) -> Result<bool> {
467        let mut stream = match connect_exec_stream(&self.socket_path).await {
468            Ok(s) => s,
469            Err(_) => return Ok(false),
470        };
471
472        let payload = format!("signal-main:{}", signal).into_bytes();
473        let frame = a3s_transport::Frame::control(payload);
474        let encoded = frame
475            .encode()
476            .map_err(|e| BoxError::ExecError(format!("signal-main frame encode failed: {}", e)))?;
477
478        if stream.write_all(&encoded).await.is_err() {
479            return Ok(false);
480        }
481
482        // Host-side deadline: a wedged guest can complete the connect handshake
483        // (listen backlog) but never write the ACK, which would hang this read
484        // forever — and stop/restart deliver the signal through here BEFORE their
485        // force-kill fallback, so the fallback would never run. On timeout report
486        // not-acknowledged so the caller force-kills.
487        let (r, _w) = tokio::io::split(stream);
488        let mut reader = a3s_transport::FrameReader::new(r);
489        let read = tokio::time::timeout(
490            std::time::Duration::from_secs(SIGNAL_MAIN_ACK_TIMEOUT_SECS),
491            reader.read_frame(),
492        )
493        .await;
494        match read {
495            Ok(Ok(Some(f)))
496                if f.frame_type == a3s_transport::FrameType::Control
497                    && f.payload == EXEC_SIGNAL_MAIN_ACK =>
498            {
499                Ok(true)
500            }
501            _ => Ok(false),
502        }
503    }
504
505    /// Ask the restricted rootfs maintenance guest to unmount its read-only
506    /// auxiliary disk and let PID 1 return. Returns false on any transport or
507    /// protocol failure so teardown can use its bounded shim fallback.
508    pub async fn shutdown_rootfs_maintenance(&self) -> Result<bool> {
509        let mut stream = match connect_exec_stream(&self.socket_path).await {
510            Ok(stream) => stream,
511            Err(_) => return Ok(false),
512        };
513        let frame = a3s_transport::Frame::control(EXEC_CONTROL_SHUTDOWN_MAINTENANCE.to_vec());
514        let encoded = frame.encode().map_err(|error| {
515            BoxError::ExecError(format!("maintenance shutdown frame encode failed: {error}"))
516        })?;
517        if stream.write_all(&encoded).await.is_err() {
518            return Ok(false);
519        }
520
521        let (read, _write) = tokio::io::split(stream);
522        let mut reader = a3s_transport::FrameReader::new(read);
523        let response = tokio::time::timeout(
524            std::time::Duration::from_secs(SIGNAL_MAIN_ACK_TIMEOUT_SECS),
525            reader.read_frame(),
526        )
527        .await;
528        Ok(matches!(
529            response,
530            Ok(Ok(Some(frame)))
531                if frame.frame_type == a3s_transport::FrameType::Control
532                    && frame.payload == EXEC_SHUTDOWN_MAINTENANCE_ACK
533        ))
534    }
535
536    /// Ask a guest that booted IDLE (`BOX_DEFERRED_MAIN=1`) to spawn its container
537    /// command — already known to the guest via BOX_EXEC_* — as the MAIN process.
538    /// The spawned main inherits the console (so its output reaches the json-file
539    /// logs) and drives the VM lifecycle. Returns `Ok(true)` if acknowledged.
540    pub async fn spawn_main(&self, spec_json: Option<&[u8]>) -> Result<bool> {
541        let mut stream = match connect_exec_stream(&self.socket_path).await {
542            Ok(s) => s,
543            Err(_) => return Ok(false),
544        };
545
546        let mut payload = b"spawn-main:".to_vec();
547        if let Some(json) = spec_json {
548            payload.extend_from_slice(json);
549        }
550        let frame = a3s_transport::Frame::control(payload);
551        let encoded = frame
552            .encode()
553            .map_err(|e| BoxError::ExecError(format!("spawn-main frame encode failed: {}", e)))?;
554
555        if stream.write_all(&encoded).await.is_err() {
556            return Ok(false);
557        }
558
559        let (r, _w) = tokio::io::split(stream);
560        let mut reader = a3s_transport::FrameReader::new(r);
561        match reader.read_frame().await {
562            Ok(Some(f))
563                if f.frame_type == a3s_transport::FrameType::Control
564                    && f.payload == EXEC_SPAWN_MAIN_ACK =>
565            {
566                Ok(true)
567            }
568            Ok(Some(f))
569                if f.frame_type == a3s_transport::FrameType::Control
570                    && f.payload.starts_with(EXEC_SPAWN_MAIN_NACK) =>
571            {
572                let reason = String::from_utf8_lossy(&f.payload[EXEC_SPAWN_MAIN_NACK.len()..]);
573                Err(BoxError::ExecError(format!(
574                    "spawn-main rejected by guest: {reason}"
575                )))
576            }
577            _ => Ok(false),
578        }
579    }
580}
581
582/// Handle for reading streaming exec events.
583///
584/// Reads frames from the exec server: Data frames contain `ExecChunk` (stdout/stderr),
585/// Control frames contain `ExecExit` (final exit code).
586pub struct StreamingExec {
587    reader: ExecFrameReader,
588    writer: Arc<Mutex<ExecFrameWriter>>,
589    started: std::time::Instant,
590    stdout_bytes: u64,
591    stderr_bytes: u64,
592    done: bool,
593}
594
595/// Cloneable input side for a running streaming exec workload.
596#[derive(Clone, Debug)]
597pub struct StreamingExecInput {
598    writer: Arc<Mutex<ExecFrameWriter>>,
599}
600
601impl StreamingExecInput {
602    /// Write bytes to the running command's stdin.
603    pub async fn write_stdin(&self, data: &[u8]) -> Result<()> {
604        self.writer
605            .lock()
606            .await
607            .write_data(data)
608            .await
609            .map_err(|e| BoxError::ExecError(format!("Streaming exec stdin write failed: {}", e)))
610    }
611
612    /// Close the running command's stdin without stopping the process.
613    pub async fn close_stdin(&self) -> Result<()> {
614        self.writer
615            .lock()
616            .await
617            .write_control(EXEC_CONTROL_STDIN_CLOSE)
618            .await
619            .map_err(|e| {
620                BoxError::ExecError(format!("Streaming exec stdin close write failed: {}", e))
621            })
622    }
623
624    /// Request cancellation of the running command.
625    pub async fn cancel(&self) -> Result<()> {
626        self.writer
627            .lock()
628            .await
629            .write_control(EXEC_CONTROL_CANCEL)
630            .await
631            .map_err(|e| BoxError::ExecError(format!("Streaming exec cancel write failed: {}", e)))
632    }
633
634    /// Deliver one of the typed Linux workload signals to the command's
635    /// process group. SIGKILL retains the established cancel control for
636    /// compatibility with older guests; SIGTERM uses the signal control.
637    pub async fn send_signal(&self, signal: ExecutionProcessSignal) -> Result<()> {
638        if signal == ExecutionProcessSignal::Kill {
639            return self.cancel().await;
640        }
641        let mut payload = EXEC_CONTROL_SIGNAL.to_vec();
642        payload.extend_from_slice(signal.linux_number().to_string().as_bytes());
643        self.writer
644            .lock()
645            .await
646            .write_control(&payload)
647            .await
648            .map_err(|e| BoxError::ExecError(format!("Streaming exec signal write failed: {}", e)))
649    }
650
651    /// Request a flush of the guest's buffered output. The guest replies with a
652    /// flush-ack (`ExecEvent::FlushAck`) once every chunk it had buffered at
653    /// flush time has been sent, establishing a clean log-rotation boundary.
654    pub async fn flush(&self) -> Result<()> {
655        self.writer
656            .lock()
657            .await
658            .write_control(EXEC_CONTROL_FLUSH)
659            .await
660            .map_err(|e| BoxError::ExecError(format!("Streaming exec flush write failed: {}", e)))
661    }
662}
663
664impl StreamingExec {
665    /// Return a cloneable input handle for this running stream.
666    pub fn input(&self) -> StreamingExecInput {
667        StreamingExecInput {
668            writer: self.writer.clone(),
669        }
670    }
671
672    /// Write bytes to the running command's stdin.
673    pub async fn write_stdin(&self, data: &[u8]) -> Result<()> {
674        self.input().write_stdin(data).await
675    }
676
677    /// Close the running command's stdin without stopping the process.
678    pub async fn close_stdin(&self) -> Result<()> {
679        self.input().close_stdin().await
680    }
681
682    /// Request a flush of the guest's buffered output (see
683    /// [`StreamingExecInput::flush`]).
684    pub async fn flush(&self) -> Result<()> {
685        self.input().flush().await
686    }
687
688    /// Read the next event from the stream.
689    ///
690    /// Returns `None` when the command has exited and all output has been read.
691    pub async fn next_event(&mut self) -> Result<Option<a3s_box_core::exec::ExecEvent>> {
692        use a3s_box_core::exec::{ExecChunk, ExecEvent, ExecExit};
693
694        if self.done {
695            return Ok(None);
696        }
697
698        let frame = match self.reader.read_frame().await {
699            Ok(Some(f)) => f,
700            Ok(None) => {
701                self.done = true;
702                return Ok(None);
703            }
704            Err(e) => {
705                self.done = true;
706                return Err(BoxError::ExecError(format!(
707                    "Streaming exec read failed: {}",
708                    e
709                )));
710            }
711        };
712
713        match frame.frame_type {
714            a3s_transport::FrameType::Data => {
715                // Data frame = ExecChunk (stdout/stderr)
716                let chunk: ExecChunk = serde_json::from_slice(&frame.payload).map_err(|e| {
717                    BoxError::ExecError(format!("Failed to parse exec chunk: {}", e))
718                })?;
719                match chunk.stream {
720                    a3s_box_core::exec::StreamType::Stdout => {
721                        self.stdout_bytes += chunk.data.len() as u64;
722                    }
723                    a3s_box_core::exec::StreamType::Stderr => {
724                        self.stderr_bytes += chunk.data.len() as u64;
725                    }
726                }
727                Ok(Some(ExecEvent::Chunk(chunk)))
728            }
729            a3s_transport::FrameType::Control => {
730                // A Control frame is either a flush-ack marker or an ExecExit.
731                if frame.payload == EXEC_FLUSH_ACK {
732                    // Boundary marker for log rotation — the stream continues.
733                    return Ok(Some(ExecEvent::FlushAck));
734                }
735                let exit: ExecExit = serde_json::from_slice(&frame.payload).map_err(|e| {
736                    BoxError::ExecError(format!("Failed to parse exec exit: {}", e))
737                })?;
738                self.done = true;
739                Ok(Some(ExecEvent::Exit(exit)))
740            }
741            a3s_transport::FrameType::Error => {
742                let msg = String::from_utf8_lossy(&frame.payload);
743                self.done = true;
744                Err(BoxError::ExecError(format!(
745                    "Streaming exec error: {}",
746                    msg
747                )))
748            }
749            _ => Err(BoxError::ExecError(format!(
750                "Unexpected frame type in stream: {:?}",
751                frame.frame_type
752            ))),
753        }
754    }
755
756    /// Request cancellation of the running streaming exec workload.
757    ///
758    /// The guest exec server treats this as a best-effort container stop signal
759    /// and should emit a final exit frame after terminating the child process.
760    pub async fn cancel(&mut self) -> Result<()> {
761        self.input().cancel().await
762    }
763
764    /// Collect all remaining output and return the final result with metrics.
765    ///
766    /// Consumes the stream, buffering all stdout/stderr until the command exits.
767    pub async fn collect(
768        mut self,
769    ) -> Result<(
770        a3s_box_core::exec::ExecOutput,
771        a3s_box_core::exec::ExecMetrics,
772    )> {
773        use a3s_box_core::exec::{ExecEvent, ExecMetrics, ExecOutput};
774
775        let mut stdout = Vec::new();
776        let mut stderr = Vec::new();
777        let mut exit_code = -1;
778        let mut truncated = false;
779
780        while let Some(event) = self.next_event().await? {
781            match event {
782                ExecEvent::Chunk(chunk) => {
783                    let target = match chunk.stream {
784                        a3s_box_core::exec::StreamType::Stdout => &mut stdout,
785                        a3s_box_core::exec::StreamType::Stderr => &mut stderr,
786                    };
787                    let remaining =
788                        a3s_box_core::exec::MAX_OUTPUT_BYTES.saturating_sub(target.len());
789                    if chunk.data.len() > remaining {
790                        truncated = true;
791                    }
792                    target.extend_from_slice(&chunk.data[..chunk.data.len().min(remaining)]);
793                }
794                ExecEvent::FlushAck => {}
795                ExecEvent::Exit(exit) => {
796                    exit_code = exit.exit_code;
797                }
798            }
799        }
800
801        let metrics = ExecMetrics {
802            duration_ms: self.started.elapsed().as_millis() as u64,
803            peak_memory_bytes: None,
804            stdout_bytes: self.stdout_bytes,
805            stderr_bytes: self.stderr_bytes,
806        };
807
808        let output = ExecOutput {
809            stdout,
810            stderr,
811            exit_code,
812            truncated,
813        };
814
815        Ok((output, metrics))
816    }
817
818    /// Whether the stream has finished (command exited or connection closed).
819    pub fn is_done(&self) -> bool {
820        self.done
821    }
822
823    /// Get execution metrics so far.
824    pub fn metrics(&self) -> a3s_box_core::exec::ExecMetrics {
825        a3s_box_core::exec::ExecMetrics {
826            duration_ms: self.started.elapsed().as_millis() as u64,
827            peak_memory_bytes: None,
828            stdout_bytes: self.stdout_bytes,
829            stderr_bytes: self.stderr_bytes,
830        }
831    }
832}
833
834impl std::fmt::Debug for StreamingExec {
835    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
836        f.debug_struct("StreamingExec")
837            .field("done", &self.done)
838            .field("stdout_bytes", &self.stdout_bytes)
839            .field("stderr_bytes", &self.stderr_bytes)
840            .finish()
841    }
842}
843
844#[cfg(all(test, unix))]
845mod tests {
846    use super::*;
847    use tokio::io::AsyncReadExt;
848    use tokio::net::UnixListener;
849
850    fn bind_test_listener(path: &Path) -> Option<UnixListener> {
851        match UnixListener::bind(path) {
852            Ok(listener) => Some(listener),
853            Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {
854                eprintln!(
855                    "skipping Unix socket test; sandbox denied bind at {}: {}",
856                    path.display(),
857                    e
858                );
859                None
860            }
861            Err(e) => panic!("failed to bind test socket {}: {}", path.display(), e),
862        }
863    }
864
865    #[tokio::test]
866    async fn test_exec_connect_nonexistent_socket() {
867        let result = ExecClient::connect(Path::new("/tmp/nonexistent-a3s-exec-test.sock")).await;
868        assert!(result.is_err());
869        let err = result.unwrap_err();
870        assert!(matches!(err, BoxError::ExecError(_)));
871    }
872
873    #[tokio::test]
874    async fn test_exec_connect_and_socket_path() {
875        let tmp = tempfile::TempDir::new().unwrap();
876        let sock_path = tmp.path().join("exec.sock");
877        let Some(_listener) = bind_test_listener(&sock_path) else {
878            return;
879        };
880
881        let client = ExecClient::connect(&sock_path).await.unwrap();
882        assert_eq!(client.socket_path(), sock_path);
883    }
884
885    #[tokio::test]
886    async fn test_exec_heartbeat_with_echo_server() {
887        let tmp = tempfile::TempDir::new().unwrap();
888        let sock_path = tmp.path().join("hb_echo.sock");
889        let Some(listener) = bind_test_listener(&sock_path) else {
890            return;
891        };
892
893        tokio::spawn(async move {
894            // Accept connect verification
895            let (stream, _) = listener.accept().await.unwrap();
896            drop(stream);
897            // Accept heartbeat connection and echo back
898            let (mut stream, _) = listener.accept().await.unwrap();
899            // Read frame header
900            let mut header = [0u8; 5];
901            stream.read_exact(&mut header).await.unwrap();
902            let len = u32::from_be_bytes([header[1], header[2], header[3], header[4]]) as usize;
903            let mut payload = vec![0u8; len];
904            if len > 0 {
905                stream.read_exact(&mut payload).await.unwrap();
906            }
907            // Respond with Heartbeat frame
908            let response = a3s_transport::Frame::heartbeat();
909            let encoded = response.encode().unwrap();
910            stream.write_all(&encoded).await.unwrap();
911        });
912
913        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
914
915        let client = ExecClient::connect(&sock_path).await.unwrap();
916        let result = client.heartbeat().await.unwrap();
917        assert!(result);
918    }
919
920    #[tokio::test]
921    async fn test_exec_heartbeat_no_response() {
922        let tmp = tempfile::TempDir::new().unwrap();
923        let sock_path = tmp.path().join("hb_close.sock");
924        let Some(listener) = bind_test_listener(&sock_path) else {
925            return;
926        };
927
928        tokio::spawn(async move {
929            // Accept connect verification
930            let (stream, _) = listener.accept().await.unwrap();
931            drop(stream);
932            // Accept heartbeat connection, read request, then close
933            let (mut stream, _) = listener.accept().await.unwrap();
934            let mut buf = vec![0u8; 1024];
935            let _ = stream.read(&mut buf).await;
936            drop(stream);
937        });
938
939        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
940
941        let client = ExecClient::connect(&sock_path).await.unwrap();
942        let result = client.heartbeat().await.unwrap();
943        assert!(!result);
944    }
945
946    #[tokio::test]
947    async fn test_exec_heartbeat_nonexistent_socket() {
948        // heartbeat() on a non-connectable socket should return false, not error
949        let client = ExecClient {
950            socket_path: PathBuf::from("/tmp/nonexistent-hb-test.sock"),
951        };
952        let result = client.heartbeat().await.unwrap();
953        assert!(!result);
954    }
955
956    #[tokio::test]
957    async fn file_transfer_uses_the_discriminated_guest_request() {
958        let tmp = tempfile::TempDir::new().unwrap();
959        let sock_path = tmp.path().join("file_transfer.sock");
960        let Some(listener) = bind_test_listener(&sock_path) else {
961            return;
962        };
963
964        tokio::spawn(async move {
965            // ExecClient::connect performs one reachability connection first.
966            let (stream, _) = listener.accept().await.unwrap();
967            drop(stream);
968            let (stream, _) = listener.accept().await.unwrap();
969            let (read, write) = tokio::io::split(stream);
970            let mut reader = a3s_transport::FrameReader::new(read);
971            let mut writer = a3s_transport::FrameWriter::new(write);
972            let frame = reader.read_frame().await.unwrap().unwrap();
973            let request: a3s_box_core::GuestSessionRequest =
974                serde_json::from_slice(&frame.payload).unwrap();
975            match request {
976                a3s_box_core::GuestSessionRequest::File(request) => {
977                    assert_eq!(request.op, a3s_box_core::FileOp::Upload);
978                    assert_eq!(request.guest_path, "~/data.bin");
979                    assert_eq!(request.data.as_deref(), Some("AAEC"));
980                    assert_eq!(request.user.as_deref(), Some("user"));
981                }
982                other => panic!("unexpected guest request: {other:?}"),
983            }
984            writer
985                .write_data(
986                    &serde_json::to_vec(&a3s_box_core::FileResponse {
987                        success: true,
988                        data: None,
989                        size: 3,
990                        error: None,
991                    })
992                    .unwrap(),
993                )
994                .await
995                .unwrap();
996        });
997
998        let client = ExecClient::connect(&sock_path).await.unwrap();
999        let response = client
1000            .file_transfer(&a3s_box_core::FileRequest {
1001                op: a3s_box_core::FileOp::Upload,
1002                guest_path: "~/data.bin".to_string(),
1003                data: Some("AAEC".to_string()),
1004                user: Some("user".to_string()),
1005                max_bytes: None,
1006            })
1007            .await
1008            .unwrap();
1009        assert!(response.success);
1010        assert_eq!(response.size, 3);
1011    }
1012
1013    #[tokio::test]
1014    async fn filesystem_uses_the_discriminated_guest_request() {
1015        let tmp = tempfile::TempDir::new().unwrap();
1016        let sock_path = tmp.path().join("filesystem.sock");
1017        let Some(listener) = bind_test_listener(&sock_path) else {
1018            return;
1019        };
1020
1021        tokio::spawn(async move {
1022            let (stream, _) = listener.accept().await.unwrap();
1023            drop(stream);
1024            let (stream, _) = listener.accept().await.unwrap();
1025            let (read, write) = tokio::io::split(stream);
1026            let mut reader = a3s_transport::FrameReader::new(read);
1027            let mut writer = a3s_transport::FrameWriter::new(write);
1028            let frame = reader.read_frame().await.unwrap().unwrap();
1029            let request: a3s_box_core::GuestSessionRequest =
1030                serde_json::from_slice(&frame.payload).unwrap();
1031            match request {
1032                a3s_box_core::GuestSessionRequest::Filesystem(request) => {
1033                    assert_eq!(request.op, a3s_box_core::FilesystemOp::ListDir);
1034                    assert_eq!(request.path, "~/data");
1035                    assert_eq!(request.depth, 2);
1036                    assert_eq!(request.user.as_deref(), Some("user"));
1037                }
1038                other => panic!("unexpected guest request: {other:?}"),
1039            }
1040            writer
1041                .write_data(
1042                    &serde_json::to_vec(&a3s_box_core::FilesystemResponse {
1043                        success: true,
1044                        entry: None,
1045                        entries: Vec::new(),
1046                        error: None,
1047                    })
1048                    .unwrap(),
1049                )
1050                .await
1051                .unwrap();
1052        });
1053
1054        let client = ExecClient::connect(&sock_path).await.unwrap();
1055        let response = client
1056            .filesystem(&a3s_box_core::FilesystemRequest {
1057                op: a3s_box_core::FilesystemOp::ListDir,
1058                path: "~/data".to_string(),
1059                destination: None,
1060                depth: 2,
1061                user: Some("user".to_string()),
1062            })
1063            .await
1064            .unwrap();
1065        assert!(response.success);
1066        assert!(response.entries.is_empty());
1067    }
1068
1069    #[tokio::test]
1070    async fn test_archive_rootfs_streams_data_until_done_marker() {
1071        let tmp = tempfile::TempDir::new().unwrap();
1072        let sock_path = tmp.path().join("archive.sock");
1073        let Some(listener) = bind_test_listener(&sock_path) else {
1074            return;
1075        };
1076
1077        tokio::spawn(async move {
1078            // ExecClient::connect performs one reachability connection first.
1079            let (stream, _) = listener.accept().await.unwrap();
1080            drop(stream);
1081            let (mut stream, _) = listener.accept().await.unwrap();
1082            let mut header = [0u8; 5];
1083            stream.read_exact(&mut header).await.unwrap();
1084            assert_eq!(header[0], a3s_transport::FrameType::Control as u8);
1085            let length = u32::from_be_bytes(header[1..5].try_into().unwrap()) as usize;
1086            let mut payload = vec![0u8; length];
1087            stream.read_exact(&mut payload).await.unwrap();
1088            assert_eq!(payload, EXEC_CONTROL_ARCHIVE_ROOTFS);
1089
1090            for payload in [b"first".as_slice(), b"-second".as_slice()] {
1091                let frame = a3s_transport::Frame::data(payload.to_vec());
1092                stream.write_all(&frame.encode().unwrap()).await.unwrap();
1093            }
1094            let done = a3s_transport::Frame::control(EXEC_ARCHIVE_ROOTFS_DONE.to_vec());
1095            stream.write_all(&done.encode().unwrap()).await.unwrap();
1096        });
1097
1098        let client = ExecClient::connect(&sock_path).await.unwrap();
1099        let output_path = tmp.path().join("rootfs.tar");
1100        let mut output = tokio::fs::File::create(&output_path).await.unwrap();
1101        let written = client.archive_rootfs(&mut output, false).await.unwrap();
1102        drop(output);
1103
1104        assert_eq!(written, 12);
1105        assert_eq!(std::fs::read(output_path).unwrap(), b"first-second");
1106    }
1107
1108    #[tokio::test]
1109    async fn test_exec_signal_main_round_trip() {
1110        let tmp = tempfile::TempDir::new().unwrap();
1111        let sock_path = tmp.path().join("signal_main.sock");
1112        let Some(listener) = bind_test_listener(&sock_path) else {
1113            return;
1114        };
1115
1116        tokio::spawn(async move {
1117            // Accept connect verification
1118            let (stream, _) = listener.accept().await.unwrap();
1119            drop(stream);
1120            // Accept signal-main connection: read the Control frame, ack it.
1121            let (stream, _) = listener.accept().await.unwrap();
1122            let (r, w) = tokio::io::split(stream);
1123            let mut reader = a3s_transport::FrameReader::new(r);
1124            let mut writer = a3s_transport::FrameWriter::new(w);
1125
1126            let frame = reader.read_frame().await.unwrap().unwrap();
1127            assert_eq!(frame.frame_type, a3s_transport::FrameType::Control);
1128            assert_eq!(frame.payload, b"signal-main:2");
1129
1130            writer.write_control(EXEC_SIGNAL_MAIN_ACK).await.unwrap();
1131        });
1132
1133        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1134
1135        let client = ExecClient::connect(&sock_path).await.unwrap();
1136        // SIGINT = 2 (image STOPSIGNAL example)
1137        let acked = client.signal_main(2).await.unwrap();
1138        assert!(acked);
1139    }
1140
1141    #[tokio::test]
1142    async fn test_rootfs_maintenance_shutdown_round_trip() {
1143        let tmp = tempfile::TempDir::new().unwrap();
1144        let sock_path = tmp.path().join("maintenance-shutdown.sock");
1145        let Some(listener) = bind_test_listener(&sock_path) else {
1146            return;
1147        };
1148
1149        tokio::spawn(async move {
1150            let (stream, _) = listener.accept().await.unwrap();
1151            drop(stream);
1152            let (stream, _) = listener.accept().await.unwrap();
1153            let (read, write) = tokio::io::split(stream);
1154            let mut reader = a3s_transport::FrameReader::new(read);
1155            let mut writer = a3s_transport::FrameWriter::new(write);
1156            let frame = reader.read_frame().await.unwrap().unwrap();
1157            assert_eq!(frame.frame_type, a3s_transport::FrameType::Control);
1158            assert_eq!(frame.payload, EXEC_CONTROL_SHUTDOWN_MAINTENANCE);
1159            writer
1160                .write_control(EXEC_SHUTDOWN_MAINTENANCE_ACK)
1161                .await
1162                .unwrap();
1163        });
1164
1165        let client = ExecClient::connect(&sock_path).await.unwrap();
1166        assert!(client.shutdown_rootfs_maintenance().await.unwrap());
1167    }
1168
1169    #[tokio::test]
1170    async fn test_exec_signal_main_nonexistent_socket() {
1171        // signal_main on a non-connectable socket returns false, not an error,
1172        // so the caller can fall back to a hard stop.
1173        let client = ExecClient {
1174            socket_path: PathBuf::from("/tmp/nonexistent-signal-main-test.sock"),
1175        };
1176        let acked = client.signal_main(15).await.unwrap();
1177        assert!(!acked);
1178    }
1179
1180    #[tokio::test]
1181    async fn test_exec_client_exec_command() {
1182        let tmp = tempfile::TempDir::new().unwrap();
1183        let sock_path = tmp.path().join("exec_cmd.sock");
1184        let Some(listener) = bind_test_listener(&sock_path) else {
1185            return;
1186        };
1187
1188        tokio::spawn(async move {
1189            // Accept connect verification
1190            let (stream, _) = listener.accept().await.unwrap();
1191            drop(stream);
1192            // Accept exec request — read Frame, respond with Frame
1193            let (stream, _) = listener.accept().await.unwrap();
1194            let (r, w) = tokio::io::split(stream);
1195            let mut reader = a3s_transport::FrameReader::new(r);
1196            let mut writer = a3s_transport::FrameWriter::new(w);
1197
1198            // Read request frame
1199            let _frame = reader.read_frame().await.unwrap().unwrap();
1200
1201            // Send response as Data frame
1202            let output = a3s_box_core::exec::ExecOutput {
1203                stdout: b"hello\n".to_vec(),
1204                stderr: vec![],
1205                exit_code: 0,
1206                truncated: false,
1207            };
1208            let payload = serde_json::to_vec(&output).unwrap();
1209            writer.write_data(&payload).await.unwrap();
1210        });
1211
1212        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1213
1214        let client = ExecClient::connect(&sock_path).await.unwrap();
1215        let req = a3s_box_core::exec::ExecRequest {
1216            request_id: None,
1217            cmd: vec!["echo".to_string(), "hello".to_string()],
1218            env: vec![],
1219            working_dir: None,
1220            rootfs: None,
1221            user: None,
1222            stdin: None,
1223            stdin_streaming: false,
1224            timeout_ns: 0,
1225            streaming: false,
1226        };
1227        let output = client.exec_command(&req).await.unwrap();
1228        assert_eq!(output.exit_code, 0);
1229        assert_eq!(&output.stdout[..], b"hello\n");
1230        assert!(output.stderr.is_empty());
1231    }
1232
1233    #[tokio::test]
1234    async fn test_exec_client_exec_stream_collect() {
1235        let tmp = tempfile::TempDir::new().unwrap();
1236        let sock_path = tmp.path().join("exec_stream.sock");
1237        let Some(listener) = bind_test_listener(&sock_path) else {
1238            return;
1239        };
1240
1241        tokio::spawn(async move {
1242            let (stream, _) = listener.accept().await.unwrap();
1243            drop(stream);
1244
1245            let (stream, _) = listener.accept().await.unwrap();
1246            let (r, w) = tokio::io::split(stream);
1247            let mut reader = a3s_transport::FrameReader::new(r);
1248            let mut writer = a3s_transport::FrameWriter::new(w);
1249
1250            let frame = reader.read_frame().await.unwrap().unwrap();
1251            let request: a3s_box_core::exec::ExecRequest =
1252                serde_json::from_slice(&frame.payload).unwrap();
1253            assert!(request.streaming);
1254
1255            let stdout = a3s_box_core::exec::ExecChunk {
1256                stream: a3s_box_core::exec::StreamType::Stdout,
1257                data: b"hello ".to_vec(),
1258            };
1259            writer
1260                .write_data(&serde_json::to_vec(&stdout).unwrap())
1261                .await
1262                .unwrap();
1263
1264            let stderr = a3s_box_core::exec::ExecChunk {
1265                stream: a3s_box_core::exec::StreamType::Stderr,
1266                data: b"warn".to_vec(),
1267            };
1268            writer
1269                .write_data(&serde_json::to_vec(&stderr).unwrap())
1270                .await
1271                .unwrap();
1272
1273            let exit = a3s_box_core::exec::ExecExit {
1274                exit_code: 17,
1275                oom_killed: false,
1276            };
1277            writer
1278                .write_control(&serde_json::to_vec(&exit).unwrap())
1279                .await
1280                .unwrap();
1281        });
1282
1283        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1284
1285        let client = ExecClient::connect(&sock_path).await.unwrap();
1286        let req = a3s_box_core::exec::ExecRequest {
1287            request_id: None,
1288            cmd: vec!["echo".to_string(), "hello".to_string()],
1289            env: vec![],
1290            working_dir: None,
1291            rootfs: None,
1292            user: None,
1293            stdin: None,
1294            stdin_streaming: false,
1295            timeout_ns: 0,
1296            streaming: false,
1297        };
1298
1299        let stream = client.exec_stream(&req).await.unwrap();
1300        let (output, metrics) = stream.collect().await.unwrap();
1301        assert_eq!(output.stdout, b"hello ");
1302        assert_eq!(output.stderr, b"warn");
1303        assert_eq!(output.exit_code, 17);
1304        assert_eq!(metrics.stdout_bytes, 6);
1305        assert_eq!(metrics.stderr_bytes, 4);
1306    }
1307
1308    #[tokio::test]
1309    async fn test_exec_client_exec_stream_cancel_writes_control_frame() {
1310        let tmp = tempfile::TempDir::new().unwrap();
1311        let sock_path = tmp.path().join("exec_stream_cancel.sock");
1312        let Some(listener) = bind_test_listener(&sock_path) else {
1313            return;
1314        };
1315
1316        tokio::spawn(async move {
1317            let (stream, _) = listener.accept().await.unwrap();
1318            drop(stream);
1319
1320            let (stream, _) = listener.accept().await.unwrap();
1321            let (r, w) = tokio::io::split(stream);
1322            let mut reader = a3s_transport::FrameReader::new(r);
1323            let mut writer = a3s_transport::FrameWriter::new(w);
1324
1325            let frame = reader.read_frame().await.unwrap().unwrap();
1326            let request: a3s_box_core::exec::ExecRequest =
1327                serde_json::from_slice(&frame.payload).unwrap();
1328            assert!(request.streaming);
1329
1330            let cancel = reader.read_frame().await.unwrap().unwrap();
1331            assert_eq!(cancel.frame_type, a3s_transport::FrameType::Control);
1332            assert_eq!(cancel.payload, b"cancel");
1333
1334            let exit = a3s_box_core::exec::ExecExit {
1335                exit_code: 137,
1336                oom_killed: false,
1337            };
1338            writer
1339                .write_control(&serde_json::to_vec(&exit).unwrap())
1340                .await
1341                .unwrap();
1342        });
1343
1344        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1345
1346        let client = ExecClient::connect(&sock_path).await.unwrap();
1347        let req = a3s_box_core::exec::ExecRequest {
1348            request_id: None,
1349            cmd: vec!["sleep".to_string(), "60".to_string()],
1350            env: vec![],
1351            working_dir: None,
1352            rootfs: None,
1353            user: None,
1354            stdin: None,
1355            stdin_streaming: false,
1356            timeout_ns: 0,
1357            streaming: false,
1358        };
1359
1360        let mut stream = client.exec_stream(&req).await.unwrap();
1361        stream.cancel().await.unwrap();
1362        let event = stream.next_event().await.unwrap().unwrap();
1363        match event {
1364            a3s_box_core::exec::ExecEvent::Exit(exit) => assert_eq!(exit.exit_code, 137),
1365            other => panic!("unexpected event: {other:?}"),
1366        }
1367    }
1368
1369    #[tokio::test]
1370    async fn test_exec_client_exec_stream_input_writes_stdin_close_and_signals() {
1371        let tmp = tempfile::TempDir::new().unwrap();
1372        let sock_path = tmp.path().join("exec_stream_stdin.sock");
1373        let Some(listener) = bind_test_listener(&sock_path) else {
1374            return;
1375        };
1376
1377        tokio::spawn(async move {
1378            let (stream, _) = listener.accept().await.unwrap();
1379            drop(stream);
1380
1381            let (stream, _) = listener.accept().await.unwrap();
1382            let (r, w) = tokio::io::split(stream);
1383            let mut reader = a3s_transport::FrameReader::new(r);
1384            let mut writer = a3s_transport::FrameWriter::new(w);
1385
1386            let frame = reader.read_frame().await.unwrap().unwrap();
1387            let request: a3s_box_core::exec::ExecRequest =
1388                serde_json::from_slice(&frame.payload).unwrap();
1389            assert!(request.streaming);
1390
1391            let stdin = reader.read_frame().await.unwrap().unwrap();
1392            assert_eq!(stdin.frame_type, a3s_transport::FrameType::Data);
1393            assert_eq!(stdin.payload, b"hello stdin\n");
1394
1395            let close = reader.read_frame().await.unwrap().unwrap();
1396            assert_eq!(close.frame_type, a3s_transport::FrameType::Control);
1397            assert_eq!(close.payload, EXEC_CONTROL_STDIN_CLOSE);
1398
1399            let terminate = reader.read_frame().await.unwrap().unwrap();
1400            assert_eq!(terminate.frame_type, a3s_transport::FrameType::Control);
1401            assert_eq!(terminate.payload, b"signal:15");
1402
1403            let kill = reader.read_frame().await.unwrap().unwrap();
1404            assert_eq!(kill.frame_type, a3s_transport::FrameType::Control);
1405            assert_eq!(kill.payload, EXEC_CONTROL_CANCEL);
1406
1407            let exit = a3s_box_core::exec::ExecExit {
1408                exit_code: 0,
1409                oom_killed: false,
1410            };
1411            writer
1412                .write_control(&serde_json::to_vec(&exit).unwrap())
1413                .await
1414                .unwrap();
1415        });
1416
1417        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1418
1419        let client = ExecClient::connect(&sock_path).await.unwrap();
1420        let req = a3s_box_core::exec::ExecRequest {
1421            request_id: None,
1422            cmd: vec!["cat".to_string()],
1423            env: vec![],
1424            working_dir: None,
1425            rootfs: None,
1426            user: None,
1427            stdin: None,
1428            stdin_streaming: true,
1429            timeout_ns: 0,
1430            streaming: false,
1431        };
1432
1433        let mut stream = client.exec_stream(&req).await.unwrap();
1434        let input = stream.input();
1435        input.write_stdin(b"hello stdin\n").await.unwrap();
1436        input.close_stdin().await.unwrap();
1437        input
1438            .send_signal(ExecutionProcessSignal::Terminate)
1439            .await
1440            .unwrap();
1441        input
1442            .send_signal(ExecutionProcessSignal::Kill)
1443            .await
1444            .unwrap();
1445        let event = stream.next_event().await.unwrap().unwrap();
1446        match event {
1447            a3s_box_core::exec::ExecEvent::Exit(exit) => assert_eq!(exit.exit_code, 0),
1448            other => panic!("unexpected event: {other:?}"),
1449        }
1450    }
1451
1452    #[tokio::test]
1453    async fn test_exec_client_flush_sends_control_and_parses_ack_then_exit() {
1454        let tmp = tempfile::TempDir::new().unwrap();
1455        let sock_path = tmp.path().join("exec_stream_flush.sock");
1456        let Some(listener) = bind_test_listener(&sock_path) else {
1457            return;
1458        };
1459
1460        tokio::spawn(async move {
1461            let (stream, _) = listener.accept().await.unwrap();
1462            drop(stream);
1463
1464            let (stream, _) = listener.accept().await.unwrap();
1465            let (r, w) = tokio::io::split(stream);
1466            let mut reader = a3s_transport::FrameReader::new(r);
1467            let mut writer = a3s_transport::FrameWriter::new(w);
1468
1469            // Consume the streaming request, then the flush control frame.
1470            let _req = reader.read_frame().await.unwrap().unwrap();
1471            let flush = reader.read_frame().await.unwrap().unwrap();
1472            assert_eq!(flush.frame_type, a3s_transport::FrameType::Control);
1473            assert_eq!(flush.payload, EXEC_CONTROL_FLUSH);
1474
1475            // Reply: a buffered chunk, the flush-ack marker, then exit.
1476            let chunk = a3s_box_core::exec::ExecChunk {
1477                stream: a3s_box_core::exec::StreamType::Stdout,
1478                data: b"pre-rotation\n".to_vec(),
1479            };
1480            writer
1481                .write_data(&serde_json::to_vec(&chunk).unwrap())
1482                .await
1483                .unwrap();
1484            writer.write_control(EXEC_FLUSH_ACK).await.unwrap();
1485            let exit = a3s_box_core::exec::ExecExit {
1486                exit_code: 0,
1487                oom_killed: false,
1488            };
1489            writer
1490                .write_control(&serde_json::to_vec(&exit).unwrap())
1491                .await
1492                .unwrap();
1493        });
1494
1495        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1496
1497        let client = ExecClient::connect(&sock_path).await.unwrap();
1498        let req = a3s_box_core::exec::ExecRequest {
1499            request_id: None,
1500            cmd: vec!["sh".to_string()],
1501            env: vec![],
1502            working_dir: None,
1503            rootfs: None,
1504            user: None,
1505            stdin: None,
1506            stdin_streaming: false,
1507            timeout_ns: 0,
1508            streaming: false,
1509        };
1510
1511        let mut stream = client.exec_stream(&req).await.unwrap();
1512        stream.flush().await.unwrap();
1513
1514        use a3s_box_core::exec::ExecEvent;
1515        match stream.next_event().await.unwrap().unwrap() {
1516            ExecEvent::Chunk(c) => assert_eq!(c.data, b"pre-rotation\n"),
1517            other => panic!("expected chunk, got {other:?}"),
1518        }
1519        // The flush-ack must parse as FlushAck, NOT as an exit (which would
1520        // wrongly end the stream).
1521        match stream.next_event().await.unwrap().unwrap() {
1522            ExecEvent::FlushAck => {}
1523            other => panic!("expected flush-ack, got {other:?}"),
1524        }
1525        match stream.next_event().await.unwrap().unwrap() {
1526            ExecEvent::Exit(exit) => assert_eq!(exit.exit_code, 0),
1527            other => panic!("expected exit, got {other:?}"),
1528        }
1529    }
1530
1531    #[tokio::test]
1532    async fn test_exec_client_malformed_response() {
1533        let tmp = tempfile::TempDir::new().unwrap();
1534        let sock_path = tmp.path().join("exec_bad.sock");
1535        let Some(listener) = bind_test_listener(&sock_path) else {
1536            return;
1537        };
1538
1539        tokio::spawn(async move {
1540            let (stream, _) = listener.accept().await.unwrap();
1541            drop(stream);
1542            let (mut stream, _) = listener.accept().await.unwrap();
1543            let mut buf = vec![0u8; 4096];
1544            let _ = stream.read(&mut buf).await;
1545            // Send garbage — not a valid frame
1546            stream.write_all(b"garbage").await.unwrap();
1547            drop(stream);
1548        });
1549
1550        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
1551
1552        let client = ExecClient::connect(&sock_path).await.unwrap();
1553        let req = a3s_box_core::exec::ExecRequest {
1554            request_id: None,
1555            cmd: vec!["test".to_string()],
1556            env: vec![],
1557            working_dir: None,
1558            rootfs: None,
1559            user: None,
1560            stdin: None,
1561            stdin_streaming: false,
1562            timeout_ns: 0,
1563            streaming: false,
1564        };
1565        let result = client.exec_command(&req).await;
1566        assert!(result.is_err());
1567    }
1568}
1569
1570#[cfg(all(test, windows))]
1571mod windows_tests {
1572    use std::sync::atomic::{AtomicU64, Ordering};
1573
1574    use tokio::net::windows::named_pipe::ServerOptions;
1575
1576    use super::*;
1577
1578    static NEXT_PIPE: AtomicU64 = AtomicU64::new(1);
1579
1580    #[tokio::test]
1581    async fn exec_command_round_trips_over_a_real_named_pipe() {
1582        let pipe_path = format!(
1583            r"\\.\pipe\a3s-box-exec-client-test-{}-{}",
1584            std::process::id(),
1585            NEXT_PIPE.fetch_add(1, Ordering::Relaxed)
1586        );
1587        let first_server = ServerOptions::new()
1588            .first_pipe_instance(true)
1589            .max_instances(254)
1590            .create(&pipe_path)
1591            .expect("create first exec pipe instance");
1592        let server_path = pipe_path.clone();
1593        let server = tokio::spawn(async move {
1594            first_server
1595                .connect()
1596                .await
1597                .expect("accept reachability connection");
1598            let second_server = ServerOptions::new()
1599                .max_instances(254)
1600                .create(&server_path)
1601                .expect("create request pipe instance");
1602            drop(first_server);
1603
1604            second_server
1605                .connect()
1606                .await
1607                .expect("accept exec request connection");
1608            let (read, write) = tokio::io::split(second_server);
1609            let mut reader = a3s_transport::FrameReader::new(read);
1610            let mut writer = a3s_transport::FrameWriter::new(write);
1611            let request = reader
1612                .read_frame()
1613                .await
1614                .expect("read request frame")
1615                .expect("request stream stays open");
1616            let request: a3s_box_core::exec::ExecRequest =
1617                serde_json::from_slice(&request.payload).expect("decode request");
1618            assert_eq!(request.cmd, ["echo", "windows"]);
1619
1620            let output = a3s_box_core::exec::ExecOutput {
1621                stdout: b"windows\n".to_vec(),
1622                stderr: Vec::new(),
1623                exit_code: 0,
1624                truncated: false,
1625            };
1626            writer
1627                .write_data(&serde_json::to_vec(&output).expect("encode response"))
1628                .await
1629                .expect("write response");
1630        });
1631
1632        let client = ExecClient::connect(Path::new(&pipe_path))
1633            .await
1634            .expect("connect exec client");
1635        let output = client
1636            .exec_command(&a3s_box_core::exec::ExecRequest {
1637                request_id: None,
1638                cmd: vec!["echo".to_string(), "windows".to_string()],
1639                timeout_ns: 5_000_000_000,
1640                env: Vec::new(),
1641                working_dir: None,
1642                rootfs: None,
1643                stdin: None,
1644                stdin_streaming: false,
1645                user: None,
1646                streaming: false,
1647            })
1648            .await
1649            .expect("execute over named pipe");
1650
1651        assert_eq!(output.stdout, b"windows\n");
1652        assert_eq!(output.exit_code, 0);
1653        server.await.expect("server task joins");
1654    }
1655}