1use 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:";
19const EXEC_CONTROL_FLUSH: &[u8] = b"flush";
21const EXEC_CONTROL_ARCHIVE_ROOTFS: &[u8] = b"archive-rootfs-v1";
23const EXEC_CONTROL_ARCHIVE_ROOTFS_PAUSE: &[u8] = b"archive-rootfs-v1:pause";
24const EXEC_ARCHIVE_ROOTFS_DONE: &[u8] = b"archive-rootfs-v1-done";
26const EXEC_FLUSH_ACK: &[u8] = b"flush-ack";
30const EXEC_SIGNAL_MAIN_ACK: &[u8] = b"signal-main-ack";
34const EXEC_SPAWN_MAIN_ACK: &[u8] = b"spawn-main-ack";
38const EXEC_SPAWN_MAIN_NACK: &[u8] = b"spawn-main-nack:";
41
42const EXEC_HOST_SLACK_SECS: u64 = 10;
46const 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#[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 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 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 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 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 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 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 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 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 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 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 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 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 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
548pub 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#[derive(Clone, Debug)]
563pub struct StreamingExecInput {
564 writer: Arc<Mutex<ExecFrameWriter>>,
565}
566
567impl StreamingExecInput {
568 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 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 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 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 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 pub fn input(&self) -> StreamingExecInput {
633 StreamingExecInput {
634 writer: self.writer.clone(),
635 }
636 }
637
638 pub async fn write_stdin(&self, data: &[u8]) -> Result<()> {
640 self.input().write_stdin(data).await
641 }
642
643 pub async fn close_stdin(&self) -> Result<()> {
645 self.input().close_stdin().await
646 }
647
648 pub async fn flush(&self) -> Result<()> {
651 self.input().flush().await
652 }
653
654 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 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 if frame.payload == EXEC_FLUSH_ACK {
698 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 pub async fn cancel(&mut self) -> Result<()> {
727 self.input().cancel().await
728 }
729
730 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 pub fn is_done(&self) -> bool {
786 self.done
787 }
788
789 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 let (stream, _) = listener.accept().await.unwrap();
862 drop(stream);
863 let (mut stream, _) = listener.accept().await.unwrap();
865 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 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 let (stream, _) = listener.accept().await.unwrap();
897 drop(stream);
898 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 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 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 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 let (stream, _) = listener.accept().await.unwrap();
1085 drop(stream);
1086 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 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 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 let (stream, _) = listener.accept().await.unwrap();
1129 drop(stream);
1130 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 let _frame = reader.read_frame().await.unwrap().unwrap();
1138
1139 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 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 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 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 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}