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 {
33 pub token: Option<Token>,
36 pub fd: Fd,
38}
39
40#[repr(align(64))]
62pub struct DrainState<F>
63where
64 F: FnMut(&[u8]) -> bool,
65{
66 pub(crate) stdout_slot: Option<FdSlot>,
67 pub(crate) stderr_slot: Option<FdSlot>,
68 pub(crate) stdin_slot: Option<FdSlot>,
69
70 pub(crate) buffer: BufferState,
71 pub(crate) writer: WriterState,
72
73 pub(crate) early_exit: Option<F>,
74
75 pub(crate) pty_master: bool,
80}
81
82impl<F> DrainState<F>
83where
84 F: FnMut(&[u8]) -> bool,
85{
86 pub fn new(
96 stdin_fd: Option<Fd>,
97 stdin_buf: Option<Box<[u8]>>,
98 stdout_fd: Option<Fd>,
99 stderr_fd: Option<Fd>,
100 limit: usize,
101 early_exit: Option<F>,
102 chunk_sink: Option<ChunkSink>,
103 pty_master: bool,
104 ) -> Result<Self, CoreError> {
105 let stdin_slot = if stdin_buf.is_some() {
106 if let Some(fd) = stdin_fd {
107 fd.set_nonblock()?;
108 Some(FdSlot { token: None, fd })
109 } else {
110 None
111 }
112 } else {
113 None
114 };
115
116 let stdout_slot = if let Some(fd) = stdout_fd {
117 fd.set_nonblock()?;
118 Some(FdSlot { token: None, fd })
119 } else {
120 None
121 };
122
123 let stderr_slot = if let Some(fd) = stderr_fd {
124 fd.set_nonblock()?;
125 Some(FdSlot { token: None, fd })
126 } else {
127 None
128 };
129
130 Ok(Self {
131 stdin_slot,
132 stdout_slot,
133 stderr_slot,
134 buffer: BufferState::new(limit, chunk_sink),
135 writer: WriterState::new(stdin_buf),
136 early_exit,
137 pty_master,
138 })
139 }
140
141 #[inline(always)]
143 pub fn is_done(&self) -> bool {
144 self.stdin_slot.is_none() && self.stdout_slot.is_none() && self.stderr_slot.is_none()
145 }
146
147 pub(crate) fn resize_pty(&self, rows: u16, cols: u16) -> Result<(), CoreError> {
158 if rows == 0 || cols == 0 {
159 return Err(CoreError::sys(
160 libc::EINVAL,
161 "resize_pty: rows and cols must be non-zero",
162 ));
163 }
164 let Some(slot) = &self.stdout_slot else {
165 return Err(CoreError::sys(libc::EINVAL, "resize_pty: no pty master"));
166 };
167 let ws = libc::winsize {
168 ws_row: rows,
169 ws_col: cols,
170 ws_xpixel: 0,
171 ws_ypixel: 0,
172 };
173 let r = unsafe { libc::ioctl(slot.fd.raw(), libc::TIOCSWINSZ as libc::Ioctl, &ws) };
174 crate::error::syscall_ret(r, "TIOCSWINSZ")
175 }
176
177 pub(crate) fn write_input(&self, bytes: &[u8]) -> Result<usize, CoreError> {
188 if !self.pty_master {
189 return Err(CoreError::sys(libc::EINVAL, "write_input: not a pty spawn"));
190 }
191 let Some(slot) = &self.stdout_slot else {
192 return Err(CoreError::sys(libc::EINVAL, "write_input: pty master closed"));
193 };
194 let fd = slot.fd.raw();
195 let mut written = 0usize;
196 while written < bytes.len() {
197 let n = unsafe {
198 libc::write(
199 fd,
200 bytes[written..].as_ptr() as *const libc::c_void,
201 bytes.len() - written,
202 )
203 };
204 if n < 0 {
205 let e = errno();
206 if e == libc::EINTR {
207 continue;
208 }
209 if e == libc::EAGAIN {
210 let mut pfd = libc::pollfd {
211 fd,
212 events: libc::POLLOUT,
213 revents: 0,
214 };
215 let rc = unsafe { libc::poll(&mut pfd, 1, WRITE_INPUT_POLL_TIMEOUT_MS) };
216 if rc < 0 {
217 let pe = errno();
218 if pe == libc::EINTR {
219 continue;
220 }
221 return Err(CoreError::sys(pe, "write_input:poll"));
222 }
223 if rc == 0 {
224 return Err(CoreError::sys(
225 libc::ETIMEDOUT,
226 "write_input: tty input buffer stayed full",
227 ));
228 }
229 continue;
230 }
231 return Err(CoreError::sys(e, "write_input"));
232 }
233 written += n as usize;
234 }
235 Ok(written)
236 }
237
238 #[inline(always)]
247 pub fn write_stdin(&mut self) -> Result<bool, CoreError> {
248 let fd = if let Some(s) = &self.stdin_slot {
249 &s.fd
250 } else {
251 return Ok(true);
252 };
253
254 let done = self.writer.write_to_fd(fd)?;
255 if done {
256 self.stdin_slot.take();
257 return Ok(true);
258 }
259 Ok(false)
260 }
261
262 #[inline(always)]
269 fn read_from_slot(
270 buffer: &mut BufferState,
271 fd: &Fd,
272 is_stdout: bool,
273 early_exit: &mut Option<F>,
274 pty_master: bool,
275 ) -> Result<ReadState, CoreError> {
276 match buffer.read_from_fd(fd, is_stdout, early_exit) {
277 Err(e) if pty_master && is_stdout && e.raw_os_error() == Some(libc::EIO) => {
278 Ok(ReadState::Eof)
279 }
280 other => other,
281 }
282 }
283
284 #[inline(always)]
294 pub fn read_fd(&mut self, is_stdout: bool) -> Result<bool, CoreError> {
295 let pty_master = self.pty_master;
296 let read_state = {
297 let slot = if is_stdout {
298 &self.stdout_slot
299 } else {
300 &self.stderr_slot
301 };
302 let fd = if let Some(s) = slot {
303 &s.fd
304 } else {
305 return Ok(true);
306 };
307 Self::read_from_slot(&mut self.buffer, fd, is_stdout, &mut self.early_exit, pty_master)?
308 };
309
310 match read_state {
311 ReadState::Open => Ok(false),
312 ReadState::Paused => Ok(false),
313 ReadState::Eof | ReadState::EarlyExit => {
314 if is_stdout {
315 self.stdout_slot.take();
316 } else {
317 self.stderr_slot.take();
318 }
319 Ok(true)
320 }
321 }
322 }
323
324 pub(crate) fn take_all_slots(&mut self) -> Vec<FdSlot> {
326 let mut slots = Vec::new();
327 if let Some(slot) = self.stdin_slot.take() {
328 slots.push(slot);
329 }
330 if let Some(slot) = self.stdout_slot.take() {
331 slots.push(slot);
332 }
333 if let Some(slot) = self.stderr_slot.take() {
334 slots.push(slot);
335 }
336 slots
337 }
338
339 pub(crate) fn register_with_reactor(
340 &mut self,
341 reactor: &mut crate::reactor::Reactor,
342 ) -> Result<(), CoreError> {
343 register_slot(reactor, &mut self.stdin_slot, false, true)?;
344 register_slot(reactor, &mut self.stdout_slot, true, false)?;
345 register_slot(reactor, &mut self.stderr_slot, true, false)?;
346 Ok(())
347 }
348
349 pub(crate) fn stdout_matches(&self, token: Token) -> bool {
350 self.stdout_slot
351 .as_ref()
352 .is_some_and(|slot| slot.token == Some(token))
353 }
354
355 pub(crate) fn stderr_matches(&self, token: Token) -> bool {
356 self.stderr_slot
357 .as_ref()
358 .is_some_and(|slot| slot.token == Some(token))
359 }
360
361 pub(crate) fn stdin_matches(&self, token: Token) -> bool {
362 self.stdin_slot
363 .as_ref()
364 .is_some_and(|slot| slot.token == Some(token))
365 }
366
367 pub(crate) fn drop_stdout(
368 &mut self,
369 reactor: &mut crate::reactor::Reactor,
370 ) -> Result<(), CoreError> {
371 if let Some(slot) = self.stdout_slot.take() {
372 del_slot(reactor, &slot)?;
373 }
374 Ok(())
375 }
376
377 pub(crate) fn drop_stderr(
378 &mut self,
379 reactor: &mut crate::reactor::Reactor,
380 ) -> Result<(), CoreError> {
381 if let Some(slot) = self.stderr_slot.take() {
382 del_slot(reactor, &slot)?;
383 }
384 Ok(())
385 }
386
387 pub(crate) fn drop_stdin(
388 &mut self,
389 reactor: &mut crate::reactor::Reactor,
390 ) -> Result<(), CoreError> {
391 if let Some(slot) = self.stdin_slot.take() {
392 del_slot(reactor, &slot)?;
393 }
394 self.writer.buf = None;
395 Ok(())
396 }
397
398 pub(crate) fn handle_stdout_ready(
399 &mut self,
400 reactor: &mut crate::reactor::Reactor,
401 ) -> Result<(), CoreError> {
402 if let Some(slot) = &self.stdout_slot {
403 let read_state = Self::read_from_slot(
404 &mut self.buffer,
405 &slot.fd,
406 true,
407 &mut self.early_exit,
408 self.pty_master,
409 )?;
410 match read_state {
411 ReadState::Open => {}
412 ReadState::Paused => {
413 self.pause_stdout(reactor)?;
418 }
419 ReadState::Eof | ReadState::EarlyExit => {
420 self.drop_stdout(reactor)?;
421 }
422 }
423 }
424 Ok(())
425 }
426
427 pub(crate) fn handle_stderr_ready(
428 &mut self,
429 reactor: &mut crate::reactor::Reactor,
430 ) -> Result<(), CoreError> {
431 if let Some(slot) = &self.stderr_slot {
432 let read_state = Self::read_from_slot(
433 &mut self.buffer,
434 &slot.fd,
435 false,
436 &mut self.early_exit,
437 self.pty_master,
438 )?;
439 match read_state {
440 ReadState::Open => {}
441 ReadState::Paused => {
442 self.pause_stderr(reactor)?;
443 }
444 ReadState::Eof | ReadState::EarlyExit => {
445 self.drop_stderr(reactor)?;
446 }
447 }
448 }
449 Ok(())
450 }
451
452 pub(crate) fn pause_stdout(
455 &mut self,
456 reactor: &mut crate::reactor::Reactor,
457 ) -> Result<(), CoreError> {
458 pause_slot(reactor, &mut self.stdout_slot)
459 }
460
461 pub(crate) fn pause_stderr(
463 &mut self,
464 reactor: &mut crate::reactor::Reactor,
465 ) -> Result<(), CoreError> {
466 pause_slot(reactor, &mut self.stderr_slot)
467 }
468
469 pub fn stdout_paused(&self) -> bool {
471 self.buffer.stdout_paused()
472 }
473
474 pub fn stderr_paused(&self) -> bool {
476 self.buffer.stderr_paused()
477 }
478
479 pub fn resume_stdout(
483 &mut self,
484 reactor: &mut crate::reactor::Reactor,
485 ) -> Result<bool, CoreError> {
486 if !self.buffer.deliver_pending_stdout()? {
487 return Ok(false);
488 }
489 register_slot(reactor, &mut self.stdout_slot, true, false)?;
490 Ok(true)
491 }
492
493 pub fn resume_stderr(
497 &mut self,
498 reactor: &mut crate::reactor::Reactor,
499 ) -> Result<bool, CoreError> {
500 if !self.buffer.deliver_pending_stderr()? {
501 return Ok(false);
502 }
503 register_slot(reactor, &mut self.stderr_slot, true, false)?;
504 Ok(true)
505 }
506
507 pub(crate) fn take_stdout_pending(&mut self) -> Option<Vec<u8>> {
509 self.buffer.take_stdout_pending()
510 }
511
512 pub(crate) fn take_stderr_pending(&mut self) -> Option<Vec<u8>> {
514 self.buffer.take_stderr_pending()
515 }
516
517 pub(crate) fn handle_stdin_writable(
518 &mut self,
519 reactor: &mut crate::reactor::Reactor,
520 ) -> Result<(), CoreError> {
521 if let Some(slot) = &self.stdin_slot {
522 let done = self.writer.write_to_fd(&slot.fd)?;
523 if done {
524 self.drop_stdin(reactor)?;
525 }
526 }
527 Ok(())
528 }
529
530 pub fn into_parts(mut self) -> (Vec<u8>, Vec<u8>) {
532 let (stdout, stderr, _, _) = std::mem::take(&mut self.buffer).into_parts();
533 (stdout, stderr)
534 }
535
536 #[inline(always)]
538 pub fn output_limit_exceeded(&self) -> bool {
539 self.buffer.output_limit_exceeded()
540 }
541
542 #[inline(always)]
544 pub fn stdout_early_exited(&self) -> bool {
545 self.buffer.stdout_early_exited()
546 }
547
548 pub(crate) fn into_parts_with_state(mut self) -> (Vec<u8>, Vec<u8>, bool, bool) {
550 std::mem::take(&mut self.buffer).into_parts()
551 }
552}
553
554fn register_slot(
557 reactor: &mut crate::reactor::Reactor,
558 slot: &mut Option<FdSlot>,
559 readable: bool,
560 writable: bool,
561) -> Result<(), CoreError> {
562 let Some(s) = slot.as_mut() else {
563 return Ok(());
564 };
565 if s.token.is_some() {
566 return Ok(());
567 }
568 s.token = Some(reactor.add(&s.fd, readable, writable)?);
569 Ok(())
570}
571
572fn del_slot(
576 reactor: &crate::reactor::Reactor,
577 slot: &FdSlot,
578) -> Result<(), CoreError> {
579 if slot.token.is_none() {
580 return Ok(());
581 }
582 match reactor.del(&slot.fd) {
583 Ok(()) => Ok(()),
584 Err(e) if e.raw_os_error() == Some(libc::ENOENT) => Ok(()),
585 Err(e) => Err(e),
586 }
587}
588
589fn pause_slot(
592 reactor: &mut crate::reactor::Reactor,
593 slot: &mut Option<FdSlot>,
594) -> Result<(), CoreError> {
595 let Some(s) = slot.as_mut() else {
596 return Ok(());
597 };
598 if s.token.is_none() {
599 return Ok(());
600 }
601 s.token = None;
602 match reactor.del(&s.fd) {
603 Ok(()) => Ok(()),
604 Err(e) if e.raw_os_error() == Some(libc::ENOENT) => Ok(()),
605 Err(e) => Err(e),
606 }
607}