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