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