1use crate::CoreError;
16use crate::fd::{Fd, Token};
17use crate::io::buffer::{BufferState, ChunkSink, ReadState};
18use crate::io::writer::WriterState;
19
20const WRITE_INPUT_POLL_TIMEOUT_MS: i32 = 2_000;
25
26#[inline(always)]
27fn errno() -> i32 {
28 std::io::Error::last_os_error().raw_os_error().unwrap_or(0)
29}
30
31pub(crate) struct FdSlot {
34 pub token: Option<Token>,
37 pub fd: Fd,
39 pub readable: bool,
41 pub writable: bool,
43}
44
45#[repr(align(64))]
67pub struct DrainState<F>
68where
69 F: FnMut(&[u8]) -> bool,
70{
71 pub(crate) stdout_slot: Option<FdSlot>,
72 pub(crate) stderr_slot: Option<FdSlot>,
73 pub(crate) stdin_slot: Option<FdSlot>,
74
75 pub(crate) buffer: BufferState,
76 pub(crate) writer: WriterState,
77
78 pub(crate) early_exit: Option<F>,
79
80 pub(crate) pty_master: bool,
85}
86
87impl<F> DrainState<F>
88where
89 F: FnMut(&[u8]) -> bool,
90{
91 #[allow(clippy::too_many_arguments)] pub fn new(
102 stdin_fd: Option<Fd>,
103 stdin_buf: Option<Box<[u8]>>,
104 stdout_fd: Option<Fd>,
105 stderr_fd: Option<Fd>,
106 limit: usize,
107 early_exit: Option<F>,
108 chunk_sink: Option<ChunkSink>,
109 pty_master: bool,
110 ) -> Result<Self, CoreError> {
111 let stdin_slot = if stdin_buf.is_some() {
112 if let Some(fd) = stdin_fd {
113 fd.set_nonblock()?;
114 Some(FdSlot {
115 token: None,
116 fd,
117 readable: false,
118 writable: false,
119 })
120 } else {
121 None
122 }
123 } else {
124 None
125 };
126
127 let stdout_slot = if let Some(fd) = stdout_fd {
128 fd.set_nonblock()?;
129 Some(FdSlot {
130 token: None,
131 fd,
132 readable: false,
133 writable: false,
134 })
135 } else {
136 None
137 };
138
139 let stderr_slot = if let Some(fd) = stderr_fd {
140 fd.set_nonblock()?;
141 Some(FdSlot {
142 token: None,
143 fd,
144 readable: false,
145 writable: false,
146 })
147 } else {
148 None
149 };
150
151 Ok(Self {
152 stdin_slot,
153 stdout_slot,
154 stderr_slot,
155 buffer: BufferState::new(limit, chunk_sink),
156 writer: WriterState::new(stdin_buf),
157 early_exit,
158 pty_master,
159 })
160 }
161
162 #[inline(always)]
164 pub fn is_done(&self) -> bool {
165 self.stdin_slot.is_none() && self.stdout_slot.is_none() && self.stderr_slot.is_none()
166 }
167
168 pub(crate) fn resize_pty(&self, rows: u16, cols: u16) -> Result<(), CoreError> {
179 if rows == 0 || cols == 0 {
180 return Err(CoreError::sys(
181 libc::EINVAL,
182 "resize_pty: rows and cols must be non-zero",
183 ));
184 }
185 let Some(slot) = &self.stdout_slot else {
186 return Err(CoreError::sys(libc::EINVAL, "resize_pty: no pty master"));
187 };
188 let ws = libc::winsize {
189 ws_row: rows,
190 ws_col: cols,
191 ws_xpixel: 0,
192 ws_ypixel: 0,
193 };
194 let r = unsafe { libc::ioctl(slot.fd.raw(), libc::TIOCSWINSZ as libc::Ioctl, &ws) };
195 crate::error::syscall_ret(r, "TIOCSWINSZ")
196 }
197
198 pub(crate) fn write_input(&self, bytes: &[u8]) -> Result<usize, CoreError> {
209 if !self.pty_master {
210 return Err(CoreError::sys(libc::EINVAL, "write_input: not a pty spawn"));
211 }
212 let Some(slot) = &self.stdout_slot else {
213 return Err(CoreError::sys(
214 libc::EINVAL,
215 "write_input: pty master closed",
216 ));
217 };
218 let fd = slot.fd.raw();
219 let mut written = 0usize;
220 while written < bytes.len() {
221 let n = unsafe {
222 libc::write(
223 fd,
224 bytes[written..].as_ptr() as *const libc::c_void,
225 bytes.len() - written,
226 )
227 };
228 if n < 0 {
229 let e = errno();
230 if e == libc::EINTR {
231 continue;
232 }
233 if e == libc::EAGAIN {
234 let mut pfd = libc::pollfd {
235 fd,
236 events: libc::POLLOUT,
237 revents: 0,
238 };
239 let rc = unsafe { libc::poll(&mut pfd, 1, WRITE_INPUT_POLL_TIMEOUT_MS) };
240 if rc < 0 {
241 let pe = errno();
242 if pe == libc::EINTR {
243 continue;
244 }
245 return Err(CoreError::sys(pe, "write_input:poll"));
246 }
247 if rc == 0 {
248 return Err(CoreError::sys(
249 libc::ETIMEDOUT,
250 "write_input: tty input buffer stayed full",
251 ));
252 }
253 continue;
254 }
255 return Err(CoreError::sys(e, "write_input"));
256 }
257 written += n as usize;
258 }
259 Ok(written)
260 }
261
262 pub(crate) fn write_input_nonblock(&self, bytes: &[u8]) -> Result<Option<usize>, CoreError> {
274 if !self.pty_master {
275 return Err(CoreError::sys(
276 libc::EINVAL,
277 "write_input_nonblock: not a pty spawn",
278 ));
279 }
280 let Some(slot) = &self.stdout_slot else {
281 return Err(CoreError::sys(
282 libc::EINVAL,
283 "write_input_nonblock: pty master closed",
284 ));
285 };
286 let fd = slot.fd.raw();
287 let mut written = 0usize;
288 while written < bytes.len() {
289 let n = unsafe {
290 libc::write(
291 fd,
292 bytes[written..].as_ptr() as *const libc::c_void,
293 bytes.len() - written,
294 )
295 };
296 if n < 0 {
297 let e = errno();
298 if e == libc::EINTR {
299 continue;
300 }
301 if e == libc::EAGAIN {
302 return Ok(if written == 0 { None } else { Some(written) });
303 }
304 return Err(CoreError::sys(e, "write_input_nonblock"));
305 }
306 written += n as usize;
307 }
308 Ok(Some(written))
309 }
310
311 #[inline(always)]
320 pub fn write_stdin(&mut self) -> Result<bool, CoreError> {
321 let fd = if let Some(s) = &self.stdin_slot {
322 &s.fd
323 } else {
324 return Ok(true);
325 };
326
327 let done = self.writer.write_to_fd(fd)?;
328 if done {
329 self.stdin_slot.take();
330 return Ok(true);
331 }
332 Ok(false)
333 }
334
335 #[inline(always)]
342 fn read_from_slot(
343 buffer: &mut BufferState,
344 fd: &Fd,
345 is_stdout: bool,
346 early_exit: &mut Option<F>,
347 pty_master: bool,
348 ) -> Result<ReadState, CoreError> {
349 match buffer.read_from_fd(fd, is_stdout, early_exit) {
350 Err(e) if pty_master && is_stdout && e.raw_os_error() == Some(libc::EIO) => {
351 Ok(ReadState::Eof)
352 }
353 other => other,
354 }
355 }
356
357 #[inline(always)]
367 pub fn read_fd(&mut self, is_stdout: bool) -> Result<bool, CoreError> {
368 let pty_master = self.pty_master;
369 let read_state = {
370 let slot = if is_stdout {
371 &self.stdout_slot
372 } else {
373 &self.stderr_slot
374 };
375 let fd = if let Some(s) = slot {
376 &s.fd
377 } else {
378 return Ok(true);
379 };
380 Self::read_from_slot(
381 &mut self.buffer,
382 fd,
383 is_stdout,
384 &mut self.early_exit,
385 pty_master,
386 )?
387 };
388
389 match read_state {
390 ReadState::Open => Ok(false),
391 ReadState::Paused => Ok(false),
392 ReadState::Eof | ReadState::EarlyExit => {
393 if is_stdout {
394 self.stdout_slot.take();
395 } else {
396 self.stderr_slot.take();
397 }
398 Ok(true)
399 }
400 }
401 }
402
403 pub(crate) fn take_all_slots(&mut self) -> Vec<FdSlot> {
405 let mut slots = Vec::new();
406 if let Some(slot) = self.stdin_slot.take() {
407 slots.push(slot);
408 }
409 if let Some(slot) = self.stdout_slot.take() {
410 slots.push(slot);
411 }
412 if let Some(slot) = self.stderr_slot.take() {
413 slots.push(slot);
414 }
415 slots
416 }
417
418 pub(crate) fn register_with_reactor(
419 &mut self,
420 reactor: &mut crate::reactor::Reactor,
421 ) -> Result<(), CoreError> {
422 register_slot(reactor, &mut self.stdin_slot, false, true)?;
423 register_slot(reactor, &mut self.stdout_slot, true, false)?;
424 register_slot(reactor, &mut self.stderr_slot, true, false)?;
425 Ok(())
426 }
427
428 pub(crate) fn stdout_matches(&self, token: Token) -> bool {
429 self.stdout_slot
430 .as_ref()
431 .is_some_and(|slot| slot.token == Some(token))
432 }
433
434 pub(crate) fn stderr_matches(&self, token: Token) -> bool {
435 self.stderr_slot
436 .as_ref()
437 .is_some_and(|slot| slot.token == Some(token))
438 }
439
440 pub(crate) fn stdin_matches(&self, token: Token) -> bool {
441 self.stdin_slot
442 .as_ref()
443 .is_some_and(|slot| slot.token == Some(token))
444 }
445
446 pub(crate) fn drop_stdout(
447 &mut self,
448 reactor: &mut crate::reactor::Reactor,
449 ) -> Result<(), CoreError> {
450 if let Some(slot) = self.stdout_slot.take() {
451 del_slot(reactor, &slot)?;
452 }
453 Ok(())
454 }
455
456 pub(crate) fn drop_stderr(
457 &mut self,
458 reactor: &mut crate::reactor::Reactor,
459 ) -> Result<(), CoreError> {
460 if let Some(slot) = self.stderr_slot.take() {
461 del_slot(reactor, &slot)?;
462 }
463 Ok(())
464 }
465
466 pub(crate) fn drop_stdin(
467 &mut self,
468 reactor: &mut crate::reactor::Reactor,
469 ) -> Result<(), CoreError> {
470 if let Some(slot) = self.stdin_slot.take() {
471 del_slot(reactor, &slot)?;
472 }
473 self.writer.buf = None;
474 Ok(())
475 }
476
477 pub(crate) fn handle_stdout_ready(
478 &mut self,
479 reactor: &mut crate::reactor::Reactor,
480 ) -> Result<(), CoreError> {
481 if let Some(slot) = &self.stdout_slot {
482 let read_state = Self::read_from_slot(
483 &mut self.buffer,
484 &slot.fd,
485 true,
486 &mut self.early_exit,
487 self.pty_master,
488 )?;
489 match read_state {
490 ReadState::Open => {}
491 ReadState::Paused => {
492 self.pause_stdout(reactor)?;
497 }
498 ReadState::Eof | ReadState::EarlyExit => {
499 self.drop_stdout(reactor)?;
500 }
501 }
502 }
503 Ok(())
504 }
505
506 pub(crate) fn handle_stderr_ready(
507 &mut self,
508 reactor: &mut crate::reactor::Reactor,
509 ) -> Result<(), CoreError> {
510 if let Some(slot) = &self.stderr_slot {
511 let read_state = Self::read_from_slot(
512 &mut self.buffer,
513 &slot.fd,
514 false,
515 &mut self.early_exit,
516 self.pty_master,
517 )?;
518 match read_state {
519 ReadState::Open => {}
520 ReadState::Paused => {
521 self.pause_stderr(reactor)?;
522 }
523 ReadState::Eof | ReadState::EarlyExit => {
524 self.drop_stderr(reactor)?;
525 }
526 }
527 }
528 Ok(())
529 }
530
531 pub(crate) fn pause_stdout(
534 &mut self,
535 reactor: &mut crate::reactor::Reactor,
536 ) -> Result<(), CoreError> {
537 pause_slot(reactor, &mut self.stdout_slot)
538 }
539
540 pub(crate) fn pause_stderr(
542 &mut self,
543 reactor: &mut crate::reactor::Reactor,
544 ) -> Result<(), CoreError> {
545 pause_slot(reactor, &mut self.stderr_slot)
546 }
547
548 pub fn stdout_paused(&self) -> bool {
550 self.buffer.stdout_paused()
551 }
552
553 pub fn stderr_paused(&self) -> bool {
555 self.buffer.stderr_paused()
556 }
557
558 pub fn resume_stdout(
562 &mut self,
563 reactor: &mut crate::reactor::Reactor,
564 ) -> Result<bool, CoreError> {
565 if !self.buffer.deliver_pending_stdout()? {
566 return Ok(false);
567 }
568 register_slot(reactor, &mut self.stdout_slot, true, false)?;
569 Ok(true)
570 }
571
572 pub fn resume_stderr(
576 &mut self,
577 reactor: &mut crate::reactor::Reactor,
578 ) -> Result<bool, CoreError> {
579 if !self.buffer.deliver_pending_stderr()? {
580 return Ok(false);
581 }
582 register_slot(reactor, &mut self.stderr_slot, true, false)?;
583 Ok(true)
584 }
585
586 pub(crate) fn take_stdout_pending(&mut self) -> Option<Vec<u8>> {
588 self.buffer.take_stdout_pending()
589 }
590
591 pub(crate) fn take_stderr_pending(&mut self) -> Option<Vec<u8>> {
593 self.buffer.take_stderr_pending()
594 }
595
596 pub(crate) fn handle_stdin_writable(
597 &mut self,
598 reactor: &mut crate::reactor::Reactor,
599 ) -> Result<(), CoreError> {
600 if let Some(slot) = &self.stdin_slot {
601 let done = self.writer.write_to_fd(&slot.fd)?;
602 if done {
603 self.drop_stdin(reactor)?;
604 }
605 }
606 Ok(())
607 }
608
609 pub fn into_parts(mut self) -> (Vec<u8>, Vec<u8>) {
611 let (stdout, stderr, _, _) = std::mem::take(&mut self.buffer).into_parts();
612 (stdout, stderr)
613 }
614
615 #[inline(always)]
617 pub fn output_limit_exceeded(&self) -> bool {
618 self.buffer.output_limit_exceeded()
619 }
620
621 #[inline(always)]
623 pub fn stdout_early_exited(&self) -> bool {
624 self.buffer.stdout_early_exited()
625 }
626
627 pub(crate) fn into_parts_with_state(mut self) -> (Vec<u8>, Vec<u8>, bool, bool) {
629 std::mem::take(&mut self.buffer).into_parts()
630 }
631}
632
633fn register_slot(
639 reactor: &mut crate::reactor::Reactor,
640 slot: &mut Option<FdSlot>,
641 readable: bool,
642 writable: bool,
643) -> Result<(), CoreError> {
644 let Some(s) = slot.as_mut() else {
645 return Ok(());
646 };
647 if let Some(token) = s.token {
648 if s.readable == readable && s.writable == writable {
649 return Ok(());
650 }
651 reactor.mod_(&s.fd, token, readable, writable)?;
652 s.readable = readable;
653 s.writable = writable;
654 return Ok(());
655 }
656 s.token = Some(reactor.add(&s.fd, readable, writable)?);
657 s.readable = readable;
658 s.writable = writable;
659 Ok(())
660}
661
662fn del_slot(reactor: &crate::reactor::Reactor, slot: &FdSlot) -> Result<(), CoreError> {
666 if slot.token.is_none() {
667 return Ok(());
668 }
669 match reactor.del(&slot.fd) {
670 Ok(()) => Ok(()),
671 Err(e) if e.raw_os_error() == Some(libc::ENOENT) => Ok(()),
672 Err(e) => Err(e),
673 }
674}
675
676fn pause_slot(
681 reactor: &mut crate::reactor::Reactor,
682 slot: &mut Option<FdSlot>,
683) -> Result<(), CoreError> {
684 let Some(s) = slot.as_mut() else {
685 return Ok(());
686 };
687 let Some(token) = s.token else {
688 return Ok(());
689 };
690 if !s.readable {
691 return Ok(());
692 }
693 match reactor.mod_(&s.fd, token, false, s.writable) {
694 Ok(()) => {
695 s.readable = false;
696 Ok(())
697 }
698 Err(e) if e.raw_os_error() == Some(libc::ENOENT) => {
699 s.token = None;
701 s.readable = false;
702 Ok(())
703 }
704 Err(e) => Err(e),
705 }
706}