std\sys\pal\windows/
pipe.rs

1use crate::ffi::OsStr;
2use crate::io::{self, BorrowedCursor, IoSlice, IoSliceMut};
3use crate::os::windows::prelude::*;
4use crate::path::Path;
5use crate::random::{DefaultRandomSource, Random};
6use crate::sync::atomic::AtomicUsize;
7use crate::sync::atomic::Ordering::Relaxed;
8use crate::sys::c;
9use crate::sys::fs::{File, OpenOptions};
10use crate::sys::handle::Handle;
11use crate::sys::pal::windows::api::{self, WinError};
12use crate::sys_common::{FromInner, IntoInner};
13use crate::{mem, ptr};
14
15////////////////////////////////////////////////////////////////////////////////
16// Anonymous pipes
17////////////////////////////////////////////////////////////////////////////////
18
19pub struct AnonPipe {
20    inner: Handle,
21}
22
23impl IntoInner<Handle> for AnonPipe {
24    fn into_inner(self) -> Handle {
25        self.inner
26    }
27}
28
29impl FromInner<Handle> for AnonPipe {
30    fn from_inner(inner: Handle) -> AnonPipe {
31        Self { inner }
32    }
33}
34
35pub struct Pipes {
36    pub ours: AnonPipe,
37    pub theirs: AnonPipe,
38}
39
40/// Although this looks similar to `anon_pipe` in the Unix module it's actually
41/// subtly different. Here we'll return two pipes in the `Pipes` return value,
42/// but one is intended for "us" where as the other is intended for "someone
43/// else".
44///
45/// Currently the only use case for this function is pipes for stdio on
46/// processes in the standard library, so "ours" is the one that'll stay in our
47/// process whereas "theirs" will be inherited to a child.
48///
49/// The ours/theirs pipes are *not* specifically readable or writable. Each
50/// one only supports a read or a write, but which is which depends on the
51/// boolean flag given. If `ours_readable` is `true`, then `ours` is readable and
52/// `theirs` is writable. Conversely, if `ours_readable` is `false`, then `ours`
53/// is writable and `theirs` is readable.
54///
55/// Also note that the `ours` pipe is always a handle opened up in overlapped
56/// mode. This means that technically speaking it should only ever be used
57/// with `OVERLAPPED` instances, but also works out ok if it's only ever used
58/// once at a time (which we do indeed guarantee).
59pub fn anon_pipe(ours_readable: bool, their_handle_inheritable: bool) -> io::Result<Pipes> {
60    // A 64kb pipe capacity is the same as a typical Linux default.
61    const PIPE_BUFFER_CAPACITY: u32 = 64 * 1024;
62
63    // Note that we specifically do *not* use `CreatePipe` here because
64    // unfortunately the anonymous pipes returned do not support overlapped
65    // operations. Instead, we create a "hopefully unique" name and create a
66    // named pipe which has overlapped operations enabled.
67    //
68    // Once we do this, we connect do it as usual via `CreateFileW`, and then
69    // we return those reader/writer halves. Note that the `ours` pipe return
70    // value is always the named pipe, whereas `theirs` is just the normal file.
71    // This should hopefully shield us from child processes which assume their
72    // stdout is a named pipe, which would indeed be odd!
73    unsafe {
74        let ours;
75        let mut name;
76        let mut tries = 0;
77        let mut reject_remote_clients_flag = c::PIPE_REJECT_REMOTE_CLIENTS;
78        loop {
79            tries += 1;
80            name = format!(
81                r"\\.\pipe\__rust_anonymous_pipe1__.{}.{}",
82                c::GetCurrentProcessId(),
83                random_number(),
84            );
85            let wide_name = OsStr::new(&name).encode_wide().chain(Some(0)).collect::<Vec<_>>();
86            let mut flags = c::FILE_FLAG_FIRST_PIPE_INSTANCE | c::FILE_FLAG_OVERLAPPED;
87            if ours_readable {
88                flags |= c::PIPE_ACCESS_INBOUND;
89            } else {
90                flags |= c::PIPE_ACCESS_OUTBOUND;
91            }
92
93            let handle = c::CreateNamedPipeW(
94                wide_name.as_ptr(),
95                flags,
96                c::PIPE_TYPE_BYTE
97                    | c::PIPE_READMODE_BYTE
98                    | c::PIPE_WAIT
99                    | reject_remote_clients_flag,
100                1,
101                PIPE_BUFFER_CAPACITY,
102                PIPE_BUFFER_CAPACITY,
103                0,
104                ptr::null_mut(),
105            );
106
107            // We pass the `FILE_FLAG_FIRST_PIPE_INSTANCE` flag above, and we're
108            // also just doing a best effort at selecting a unique name. If
109            // `ERROR_ACCESS_DENIED` is returned then it could mean that we
110            // accidentally conflicted with an already existing pipe, so we try
111            // again.
112            //
113            // Don't try again too much though as this could also perhaps be a
114            // legit error.
115            // If `ERROR_INVALID_PARAMETER` is returned, this probably means we're
116            // running on pre-Vista version where `PIPE_REJECT_REMOTE_CLIENTS` is
117            // not supported, so we continue retrying without it. This implies
118            // reduced security on Windows versions older than Vista by allowing
119            // connections to this pipe from remote machines.
120            // Proper fix would increase the number of FFI imports and introduce
121            // significant amount of Windows XP specific code with no clean
122            // testing strategy
123            // For more info, see https://github.com/rust-lang/rust/pull/37677.
124            if handle == c::INVALID_HANDLE_VALUE {
125                let error = api::get_last_error();
126                if tries < 10 {
127                    if error == WinError::ACCESS_DENIED {
128                        continue;
129                    } else if reject_remote_clients_flag != 0
130                        && error == WinError::INVALID_PARAMETER
131                    {
132                        reject_remote_clients_flag = 0;
133                        tries -= 1;
134                        continue;
135                    }
136                }
137                return Err(io::Error::from_raw_os_error(error.code as i32));
138            }
139            ours = Handle::from_raw_handle(handle);
140            break;
141        }
142
143        // Connect to the named pipe we just created. This handle is going to be
144        // returned in `theirs`, so if `ours` is readable we want this to be
145        // writable, otherwise if `ours` is writable we want this to be
146        // readable.
147        //
148        // Additionally we don't enable overlapped mode on this because most
149        // client processes aren't enabled to work with that.
150        let mut opts = OpenOptions::new();
151        opts.write(ours_readable);
152        opts.read(!ours_readable);
153        opts.share_mode(0);
154        let size = mem::size_of::<c::SECURITY_ATTRIBUTES>();
155        let mut sa = c::SECURITY_ATTRIBUTES {
156            nLength: size as u32,
157            lpSecurityDescriptor: ptr::null_mut(),
158            bInheritHandle: their_handle_inheritable as i32,
159        };
160        opts.security_attributes(&mut sa);
161        let theirs = File::open(Path::new(&name), &opts)?;
162        let theirs = AnonPipe { inner: theirs.into_inner() };
163
164        Ok(Pipes {
165            ours: AnonPipe { inner: ours },
166            theirs: AnonPipe { inner: theirs.into_inner() },
167        })
168    }
169}
170
171/// Takes an asynchronous source pipe and returns a synchronous pipe suitable
172/// for sending to a child process.
173///
174/// This is achieved by creating a new set of pipes and spawning a thread that
175/// relays messages between the source and the synchronous pipe.
176pub fn spawn_pipe_relay(
177    source: &AnonPipe,
178    ours_readable: bool,
179    their_handle_inheritable: bool,
180) -> io::Result<AnonPipe> {
181    // We need this handle to live for the lifetime of the thread spawned below.
182    let source = source.try_clone()?;
183
184    // create a new pair of anon pipes.
185    let Pipes { theirs, ours } = anon_pipe(ours_readable, their_handle_inheritable)?;
186
187    // Spawn a thread that passes messages from one pipe to the other.
188    // Any errors will simply cause the thread to exit.
189    let (reader, writer) = if ours_readable { (ours, source) } else { (source, ours) };
190    crate::thread::spawn(move || {
191        let mut buf = [0_u8; 4096];
192        'reader: while let Ok(len) = reader.read(&mut buf) {
193            if len == 0 {
194                break;
195            }
196            let mut start = 0;
197            while let Ok(written) = writer.write(&buf[start..len]) {
198                start += written;
199                if start == len {
200                    continue 'reader;
201                }
202            }
203            break;
204        }
205    });
206
207    // Return the pipe that should be sent to the child process.
208    Ok(theirs)
209}
210
211fn random_number() -> usize {
212    static N: AtomicUsize = AtomicUsize::new(0);
213    loop {
214        if N.load(Relaxed) != 0 {
215            return N.fetch_add(1, Relaxed);
216        }
217
218        N.store(usize::random(&mut DefaultRandomSource), Relaxed);
219    }
220}
221
222impl AnonPipe {
223    pub fn handle(&self) -> &Handle {
224        &self.inner
225    }
226    pub fn into_handle(self) -> Handle {
227        self.inner
228    }
229
230    pub fn try_clone(&self) -> io::Result<Self> {
231        self.inner.duplicate(0, false, c::DUPLICATE_SAME_ACCESS).map(|inner| AnonPipe { inner })
232    }
233
234    pub fn read(&self, buf: &mut [u8]) -> io::Result<usize> {
235        let result = unsafe {
236            let len = crate::cmp::min(buf.len(), u32::MAX as usize) as u32;
237            let ptr = buf.as_mut_ptr();
238            self.alertable_io_internal(|overlapped, callback| {
239                c::ReadFileEx(self.inner.as_raw_handle(), ptr, len, overlapped, callback)
240            })
241        };
242
243        match result {
244            // The special treatment of BrokenPipe is to deal with Windows
245            // pipe semantics, which yields this error when *reading* from
246            // a pipe after the other end has closed; we interpret that as
247            // EOF on the pipe.
248            Err(ref e) if e.kind() == io::ErrorKind::BrokenPipe => Ok(0),
249            _ => result,
250        }
251    }
252
253    pub fn read_buf(&self, mut buf: BorrowedCursor<'_>) -> io::Result<()> {
254        let result = unsafe {
255            let len = crate::cmp::min(buf.capacity(), u32::MAX as usize) as u32;
256            let ptr = buf.as_mut().as_mut_ptr().cast::<u8>();
257            self.alertable_io_internal(|overlapped, callback| {
258                c::ReadFileEx(self.inner.as_raw_handle(), ptr, len, overlapped, callback)
259            })
260        };
261
262        match result {
263            // The special treatment of BrokenPipe is to deal with Windows
264            // pipe semantics, which yields this error when *reading* from
265            // a pipe after the other end has closed; we interpret that as
266            // EOF on the pipe.
267            Err(ref e) if e.kind() == io::ErrorKind::BrokenPipe => Ok(()),
268            Err(e) => Err(e),
269            Ok(n) => {
270                unsafe {
271                    buf.advance_unchecked(n);
272                }
273                Ok(())
274            }
275        }
276    }
277
278    pub fn read_vectored(&self, bufs: &mut [IoSliceMut<'_>]) -> io::Result<usize> {
279        self.inner.read_vectored(bufs)
280    }
281
282    #[inline]
283    pub fn is_read_vectored(&self) -> bool {
284        self.inner.is_read_vectored()
285    }
286
287    pub fn read_to_end(&self, buf: &mut Vec<u8>) -> io::Result<usize> {
288        self.handle().read_to_end(buf)
289    }
290
291    pub fn write(&self, buf: &[u8]) -> io::Result<usize> {
292        unsafe {
293            let len = crate::cmp::min(buf.len(), u32::MAX as usize) as u32;
294            self.alertable_io_internal(|overlapped, callback| {
295                c::WriteFileEx(self.inner.as_raw_handle(), buf.as_ptr(), len, overlapped, callback)
296            })
297        }
298    }
299
300    pub fn write_vectored(&self, bufs: &[IoSlice<'_>]) -> io::Result<usize> {
301        self.inner.write_vectored(bufs)
302    }
303
304    #[inline]
305    pub fn is_write_vectored(&self) -> bool {
306        self.inner.is_write_vectored()
307    }
308
309    /// Synchronizes asynchronous reads or writes using our anonymous pipe.
310    ///
311    /// This is a wrapper around [`ReadFileEx`] or [`WriteFileEx`] that uses
312    /// [Asynchronous Procedure Call] (APC) to synchronize reads or writes.
313    ///
314    /// Note: This should not be used for handles we don't create.
315    ///
316    /// # Safety
317    ///
318    /// `buf` must be a pointer to a buffer that's valid for reads or writes
319    /// up to `len` bytes. The `AlertableIoFn` must be either `ReadFileEx` or `WriteFileEx`
320    ///
321    /// [`ReadFileEx`]: https://docs.microsoft.com/en-us/windows/win32/api/fileapi/nf-fileapi-readfileex
322    /// [`WriteFileEx`]: https://docs.microsoft.com/en-us/windows/win32/api/fileapi/nf-fileapi-writefileex
323    /// [Asynchronous Procedure Call]: https://docs.microsoft.com/en-us/windows/win32/sync/asynchronous-procedure-calls
324    unsafe fn alertable_io_internal(
325        &self,
326        io: impl FnOnce(&mut c::OVERLAPPED, c::LPOVERLAPPED_COMPLETION_ROUTINE) -> c::BOOL,
327    ) -> io::Result<usize> {
328        // Use "alertable I/O" to synchronize the pipe I/O.
329        // This has four steps.
330        //
331        // STEP 1: Start the asynchronous I/O operation.
332        //         This simply calls either `ReadFileEx` or `WriteFileEx`,
333        //         giving it a pointer to the buffer and callback function.
334        //
335        // STEP 2: Enter an alertable state.
336        //         The callback set in step 1 will not be called until the thread
337        //         enters an "alertable" state. This can be done using `SleepEx`.
338        //
339        // STEP 3: The callback
340        //         Once the I/O is complete and the thread is in an alertable state,
341        //         the callback will be run on the same thread as the call to
342        //         `ReadFileEx` or `WriteFileEx` done in step 1.
343        //         In the callback we simply set the result of the async operation.
344        //
345        // STEP 4: Return the result.
346        //         At this point we'll have a result from the callback function
347        //         and can simply return it. Note that we must not return earlier,
348        //         while the I/O is still in progress.
349
350        // The result that will be set from the asynchronous callback.
351        let mut async_result: Option<AsyncResult> = None;
352        struct AsyncResult {
353            error: u32,
354            transferred: u32,
355        }
356
357        // STEP 3: The callback.
358        unsafe extern "system" fn callback(
359            dwErrorCode: u32,
360            dwNumberOfBytesTransferred: u32,
361            lpOverlapped: *mut c::OVERLAPPED,
362        ) {
363            // Set `async_result` using a pointer smuggled through `hEvent`.
364            // SAFETY:
365            // At this point, the OVERLAPPED struct will have been written to by the OS,
366            // except for our `hEvent` field which we set to a valid AsyncResult pointer (see below)
367            unsafe {
368                let result =
369                    AsyncResult { error: dwErrorCode, transferred: dwNumberOfBytesTransferred };
370                *(*lpOverlapped).hEvent.cast::<Option<AsyncResult>>() = Some(result);
371            }
372        }
373
374        // STEP 1: Start the I/O operation.
375        let mut overlapped: c::OVERLAPPED = unsafe { crate::mem::zeroed() };
376        // `hEvent` is unused by `ReadFileEx` and `WriteFileEx`.
377        // Therefore the documentation suggests using it to smuggle a pointer to the callback.
378        overlapped.hEvent = (&raw mut async_result) as *mut _;
379
380        // Asynchronous read of the pipe.
381        // If successful, `callback` will be called once it completes.
382        let result = io(&mut overlapped, Some(callback));
383        if result == c::FALSE {
384            // We can return here because the call failed.
385            // After this we must not return until the I/O completes.
386            return Err(io::Error::last_os_error());
387        }
388
389        // Wait indefinitely for the result.
390        let result = loop {
391            // STEP 2: Enter an alertable state.
392            // The second parameter of `SleepEx` is used to make this sleep alertable.
393            unsafe { c::SleepEx(c::INFINITE, c::TRUE) };
394            if let Some(result) = async_result {
395                break result;
396            }
397        };
398        // STEP 4: Return the result.
399        // `async_result` is always `Some` at this point
400        match result.error {
401            c::ERROR_SUCCESS => Ok(result.transferred as usize),
402            error => Err(io::Error::from_raw_os_error(error as _)),
403        }
404    }
405}
406
407pub fn read2(p1: AnonPipe, v1: &mut Vec<u8>, p2: AnonPipe, v2: &mut Vec<u8>) -> io::Result<()> {
408    let p1 = p1.into_handle();
409    let p2 = p2.into_handle();
410
411    let mut p1 = AsyncPipe::new(p1, v1)?;
412    let mut p2 = AsyncPipe::new(p2, v2)?;
413    let objs = [p1.event.as_raw_handle(), p2.event.as_raw_handle()];
414
415    // In a loop we wait for either pipe's scheduled read operation to complete.
416    // If the operation completes with 0 bytes, that means EOF was reached, in
417    // which case we just finish out the other pipe entirely.
418    //
419    // Note that overlapped I/O is in general super unsafe because we have to
420    // be careful to ensure that all pointers in play are valid for the entire
421    // duration of the I/O operation (where tons of operations can also fail).
422    // The destructor for `AsyncPipe` ends up taking care of most of this.
423    loop {
424        let res = unsafe { c::WaitForMultipleObjects(2, objs.as_ptr(), c::FALSE, c::INFINITE) };
425        if res == c::WAIT_OBJECT_0 {
426            if !p1.result()? || !p1.schedule_read()? {
427                return p2.finish();
428            }
429        } else if res == c::WAIT_OBJECT_0 + 1 {
430            if !p2.result()? || !p2.schedule_read()? {
431                return p1.finish();
432            }
433        } else {
434            return Err(io::Error::last_os_error());
435        }
436    }
437}
438
439struct AsyncPipe<'a> {
440    pipe: Handle,
441    event: Handle,
442    overlapped: Box<c::OVERLAPPED>, // needs a stable address
443    dst: &'a mut Vec<u8>,
444    state: State,
445}
446
447#[derive(PartialEq, Debug)]
448enum State {
449    NotReading,
450    Reading,
451    Read(usize),
452}
453
454impl<'a> AsyncPipe<'a> {
455    fn new(pipe: Handle, dst: &'a mut Vec<u8>) -> io::Result<AsyncPipe<'a>> {
456        // Create an event which we'll use to coordinate our overlapped
457        // operations, this event will be used in WaitForMultipleObjects
458        // and passed as part of the OVERLAPPED handle.
459        //
460        // Note that we do a somewhat clever thing here by flagging the
461        // event as being manually reset and setting it initially to the
462        // signaled state. This means that we'll naturally fall through the
463        // WaitForMultipleObjects call above for pipes created initially,
464        // and the only time an even will go back to "unset" will be once an
465        // I/O operation is successfully scheduled (what we want).
466        let event = Handle::new_event(true, true)?;
467        let mut overlapped: Box<c::OVERLAPPED> = unsafe { Box::new(mem::zeroed()) };
468        overlapped.hEvent = event.as_raw_handle();
469        Ok(AsyncPipe { pipe, overlapped, event, dst, state: State::NotReading })
470    }
471
472    /// Executes an overlapped read operation.
473    ///
474    /// Must not currently be reading, and returns whether the pipe is currently
475    /// at EOF or not. If the pipe is not at EOF then `result()` must be called
476    /// to complete the read later on (may block), but if the pipe is at EOF
477    /// then `result()` should not be called as it will just block forever.
478    fn schedule_read(&mut self) -> io::Result<bool> {
479        assert_eq!(self.state, State::NotReading);
480        let amt = unsafe {
481            if self.dst.capacity() == self.dst.len() {
482                let additional = if self.dst.capacity() == 0 { 16 } else { 1 };
483                self.dst.reserve(additional);
484            }
485            self.pipe.read_overlapped(self.dst.spare_capacity_mut(), &mut *self.overlapped)?
486        };
487
488        // If this read finished immediately then our overlapped event will
489        // remain signaled (it was signaled coming in here) and we'll progress
490        // down to the method below.
491        //
492        // Otherwise the I/O operation is scheduled and the system set our event
493        // to not signaled, so we flag ourselves into the reading state and move
494        // on.
495        self.state = match amt {
496            Some(0) => return Ok(false),
497            Some(amt) => State::Read(amt),
498            None => State::Reading,
499        };
500        Ok(true)
501    }
502
503    /// Wait for the result of the overlapped operation previously executed.
504    ///
505    /// Takes a parameter `wait` which indicates if this pipe is currently being
506    /// read whether the function should block waiting for the read to complete.
507    ///
508    /// Returns values:
509    ///
510    /// * `true` - finished any pending read and the pipe is not at EOF (keep
511    ///            going)
512    /// * `false` - finished any pending read and pipe is at EOF (stop issuing
513    ///             reads)
514    fn result(&mut self) -> io::Result<bool> {
515        let amt = match self.state {
516            State::NotReading => return Ok(true),
517            State::Reading => self.pipe.overlapped_result(&mut *self.overlapped, true)?,
518            State::Read(amt) => amt,
519        };
520        self.state = State::NotReading;
521        unsafe {
522            let len = self.dst.len();
523            self.dst.set_len(len + amt);
524        }
525        Ok(amt != 0)
526    }
527
528    /// Finishes out reading this pipe entirely.
529    ///
530    /// Waits for any pending and schedule read, and then calls `read_to_end`
531    /// if necessary to read all the remaining information.
532    fn finish(&mut self) -> io::Result<()> {
533        while self.result()? && self.schedule_read()? {
534            // ...
535        }
536        Ok(())
537    }
538}
539
540impl<'a> Drop for AsyncPipe<'a> {
541    fn drop(&mut self) {
542        match self.state {
543            State::Reading => {}
544            _ => return,
545        }
546
547        // If we have a pending read operation, then we have to make sure that
548        // it's *done* before we actually drop this type. The kernel requires
549        // that the `OVERLAPPED` and buffer pointers are valid for the entire
550        // I/O operation.
551        //
552        // To do that, we call `CancelIo` to cancel any pending operation, and
553        // if that succeeds we wait for the overlapped result.
554        //
555        // If anything here fails, there's not really much we can do, so we leak
556        // the buffer/OVERLAPPED pointers to ensure we're at least memory safe.
557        if self.pipe.cancel_io().is_err() || self.result().is_err() {
558            let buf = mem::take(self.dst);
559            let overlapped = Box::new(unsafe { mem::zeroed() });
560            let overlapped = mem::replace(&mut self.overlapped, overlapped);
561            mem::forget((buf, overlapped));
562        }
563    }
564}