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