coreshift_core/io/
drain.rs1use crate::CoreError;
16use crate::fd::{Fd, Token};
17use crate::io::buffer::{BufferState, ReadState};
18use crate::io::writer::WriterState;
19
20pub(crate) struct FdSlot {
22 pub token: Option<Token>,
24 pub fd: Fd,
26}
27
28#[repr(align(64))]
50pub struct DrainState<F>
51where
52 F: FnMut(&[u8]) -> bool,
53{
54 pub(crate) stdout_slot: Option<FdSlot>,
55 pub(crate) stderr_slot: Option<FdSlot>,
56 pub(crate) stdin_slot: Option<FdSlot>,
57
58 pub(crate) buffer: BufferState,
59 pub(crate) writer: WriterState,
60
61 pub(crate) early_exit: Option<F>,
62}
63
64impl<F> DrainState<F>
65where
66 F: FnMut(&[u8]) -> bool,
67{
68 pub fn new(
75 stdin_fd: Option<Fd>,
76 stdin_buf: Option<Box<[u8]>>,
77 stdout_fd: Option<Fd>,
78 stderr_fd: Option<Fd>,
79 limit: usize,
80 early_exit: Option<F>,
81 ) -> Result<Self, CoreError> {
82 let stdin_slot = if stdin_buf.is_some() {
83 if let Some(fd) = stdin_fd {
84 fd.set_nonblock()?;
85 Some(FdSlot { token: None, fd })
86 } else {
87 None
88 }
89 } else {
90 None
91 };
92
93 let stdout_slot = if let Some(fd) = stdout_fd {
94 fd.set_nonblock()?;
95 Some(FdSlot { token: None, fd })
96 } else {
97 None
98 };
99
100 let stderr_slot = if let Some(fd) = stderr_fd {
101 fd.set_nonblock()?;
102 Some(FdSlot { token: None, fd })
103 } else {
104 None
105 };
106
107 Ok(Self {
108 stdin_slot,
109 stdout_slot,
110 stderr_slot,
111 buffer: BufferState::new(limit),
112 writer: WriterState::new(stdin_buf),
113 early_exit,
114 })
115 }
116
117 #[inline(always)]
119 pub fn is_done(&self) -> bool {
120 self.stdin_slot.is_none() && self.stdout_slot.is_none() && self.stderr_slot.is_none()
121 }
122
123 #[inline(always)]
132 pub fn write_stdin(&mut self) -> Result<bool, CoreError> {
133 let fd = if let Some(s) = &self.stdin_slot {
134 &s.fd
135 } else {
136 return Ok(true);
137 };
138
139 let done = self.writer.write_to_fd(fd)?;
140 if done {
141 self.stdin_slot.take();
142 return Ok(true);
143 }
144 Ok(false)
145 }
146
147 #[inline(always)]
156 pub fn read_fd(&mut self, is_stdout: bool) -> Result<bool, CoreError> {
157 let read_state = {
158 let slot = if is_stdout {
159 &self.stdout_slot
160 } else {
161 &self.stderr_slot
162 };
163 let fd = if let Some(s) = slot {
164 &s.fd
165 } else {
166 return Ok(true);
167 };
168 self.buffer
169 .read_from_fd(fd, is_stdout, &mut self.early_exit)?
170 };
171
172 if read_state != ReadState::Open {
173 if is_stdout {
174 self.stdout_slot.take();
175 return Ok(true);
176 } else {
177 self.stderr_slot.take();
178 return Ok(true);
179 }
180 }
181
182 Ok(false)
183 }
184
185 pub(crate) fn take_all_slots(&mut self) -> Vec<FdSlot> {
187 let mut slots = Vec::new();
188 if let Some(slot) = self.stdin_slot.take() {
189 slots.push(slot);
190 }
191 if let Some(slot) = self.stdout_slot.take() {
192 slots.push(slot);
193 }
194 if let Some(slot) = self.stderr_slot.take() {
195 slots.push(slot);
196 }
197 slots
198 }
199
200 pub(crate) fn register_with_reactor(
201 &mut self,
202 reactor: &mut crate::reactor::Reactor,
203 ) -> Result<(), CoreError> {
204 register_slot(reactor, &mut self.stdin_slot, false, true)?;
205 register_slot(reactor, &mut self.stdout_slot, true, false)?;
206 register_slot(reactor, &mut self.stderr_slot, true, false)?;
207 Ok(())
208 }
209
210 pub(crate) fn stdout_matches(&self, token: Token) -> bool {
211 self.stdout_slot
212 .as_ref()
213 .is_some_and(|slot| slot.token == Some(token))
214 }
215
216 pub(crate) fn stderr_matches(&self, token: Token) -> bool {
217 self.stderr_slot
218 .as_ref()
219 .is_some_and(|slot| slot.token == Some(token))
220 }
221
222 pub(crate) fn stdin_matches(&self, token: Token) -> bool {
223 self.stdin_slot
224 .as_ref()
225 .is_some_and(|slot| slot.token == Some(token))
226 }
227
228 pub(crate) fn drop_stdout(
229 &mut self,
230 reactor: &mut crate::reactor::Reactor,
231 ) -> Result<(), CoreError> {
232 if let Some(slot) = self.stdout_slot.take() {
233 reactor.del(&slot.fd)?;
234 }
235 Ok(())
236 }
237
238 pub(crate) fn drop_stderr(
239 &mut self,
240 reactor: &mut crate::reactor::Reactor,
241 ) -> Result<(), CoreError> {
242 if let Some(slot) = self.stderr_slot.take() {
243 reactor.del(&slot.fd)?;
244 }
245 Ok(())
246 }
247
248 pub(crate) fn drop_stdin(
249 &mut self,
250 reactor: &mut crate::reactor::Reactor,
251 ) -> Result<(), CoreError> {
252 if let Some(slot) = self.stdin_slot.take() {
253 reactor.del(&slot.fd)?;
254 }
255 self.writer.buf = None;
256 Ok(())
257 }
258
259 pub(crate) fn handle_stdout_ready(
260 &mut self,
261 reactor: &mut crate::reactor::Reactor,
262 ) -> Result<(), CoreError> {
263 if let Some(slot) = &self.stdout_slot {
264 let read_state = self
265 .buffer
266 .read_from_fd(&slot.fd, true, &mut self.early_exit)?;
267 if read_state != ReadState::Open {
268 self.drop_stdout(reactor)?;
269 }
270 }
271 Ok(())
272 }
273
274 pub(crate) fn handle_stderr_ready(
275 &mut self,
276 reactor: &mut crate::reactor::Reactor,
277 ) -> Result<(), CoreError> {
278 if let Some(slot) = &self.stderr_slot {
279 let read_state = self
280 .buffer
281 .read_from_fd(&slot.fd, false, &mut self.early_exit)?;
282 if read_state != ReadState::Open {
283 self.drop_stderr(reactor)?;
284 }
285 }
286 Ok(())
287 }
288
289 pub(crate) fn handle_stdin_writable(
290 &mut self,
291 reactor: &mut crate::reactor::Reactor,
292 ) -> Result<(), CoreError> {
293 if let Some(slot) = &self.stdin_slot {
294 let done = self.writer.write_to_fd(&slot.fd)?;
295 if done {
296 self.drop_stdin(reactor)?;
297 }
298 }
299 Ok(())
300 }
301
302 pub fn into_parts(mut self) -> (Vec<u8>, Vec<u8>) {
304 let (stdout, stderr, _, _) = std::mem::take(&mut self.buffer).into_parts();
305 (stdout, stderr)
306 }
307
308 #[inline(always)]
310 pub fn output_limit_exceeded(&self) -> bool {
311 self.buffer.output_limit_exceeded()
312 }
313
314 #[inline(always)]
316 pub fn stdout_early_exited(&self) -> bool {
317 self.buffer.stdout_early_exited()
318 }
319
320 pub(crate) fn into_parts_with_state(mut self) -> (Vec<u8>, Vec<u8>, bool, bool) {
322 std::mem::take(&mut self.buffer).into_parts()
323 }
324}
325
326fn register_slot(
329 reactor: &mut crate::reactor::Reactor,
330 slot: &mut Option<FdSlot>,
331 readable: bool,
332 writable: bool,
333) -> Result<(), CoreError> {
334 let Some(s) = slot.as_mut() else {
335 return Ok(());
336 };
337 if s.token.is_some() {
338 return Ok(());
339 }
340 s.token = Some(reactor.add(&s.fd, readable, writable)?);
341 Ok(())
342}