|
| 1 | +// Copyright (c) 2018 The rust-gpio-cdev Project Developers. |
| 2 | +// |
| 3 | +// Licensed under the Apache License, Version 2.0 <LICENSE-APACHE or |
| 4 | +// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license |
| 5 | +// <LICENSE-MIT or http://opensource.org/licenses/MIT>, at your |
| 6 | +// option. This file may not be copied, modified, or distributed |
| 7 | +// except according to those terms. |
| 8 | + |
| 9 | +//! Wrapper for asynchronous programming using Tokio. |
| 10 | +
|
| 11 | +use futures::ready; |
| 12 | +use futures::stream::Stream; |
| 13 | +use futures::task::{Context, Poll}; |
| 14 | +use mio::event::Evented; |
| 15 | +use mio::unix::EventedFd; |
| 16 | +use mio::{PollOpt, Ready, Token}; |
| 17 | +use tokio::io::PollEvented; |
| 18 | + |
| 19 | +use std::io; |
| 20 | +use std::mem; |
| 21 | +use std::os::unix::io::AsRawFd; |
| 22 | +use std::pin::Pin; |
| 23 | +use std::slice; |
| 24 | + |
| 25 | +use super::errors::event_err; |
| 26 | +use super::{ffi, LineEvent, LineEventHandle, Result}; |
| 27 | + |
| 28 | +struct PollWrapper { |
| 29 | + handle: LineEventHandle, |
| 30 | +} |
| 31 | + |
| 32 | +impl Evented for PollWrapper { |
| 33 | + fn register( |
| 34 | + &self, |
| 35 | + poll: &mio::Poll, |
| 36 | + token: Token, |
| 37 | + interest: Ready, |
| 38 | + opts: PollOpt, |
| 39 | + ) -> io::Result<()> { |
| 40 | + EventedFd(&self.handle.file.as_raw_fd()).register(poll, token, interest, opts) |
| 41 | + } |
| 42 | + |
| 43 | + fn reregister( |
| 44 | + &self, |
| 45 | + poll: &mio::Poll, |
| 46 | + token: Token, |
| 47 | + interest: Ready, |
| 48 | + opts: PollOpt, |
| 49 | + ) -> io::Result<()> { |
| 50 | + EventedFd(&self.handle.file.as_raw_fd()).reregister(poll, token, interest, opts) |
| 51 | + } |
| 52 | + |
| 53 | + fn deregister(&self, poll: &mio::Poll) -> io::Result<()> { |
| 54 | + EventedFd(&self.handle.file.as_raw_fd()).deregister(poll) |
| 55 | + } |
| 56 | +} |
| 57 | + |
| 58 | +/// Wrapper around a `LineEventHandle` which implements a `futures::stream::Stream` for interrupts. |
| 59 | +/// |
| 60 | +/// # Example |
| 61 | +/// |
| 62 | +/// The following example waits for state changes on an input line. |
| 63 | +/// |
| 64 | +/// ```no_run |
| 65 | +/// # type Result<T> = std::result::Result<T, gpio_cdev::errors::Error>; |
| 66 | +/// use futures::stream::StreamExt; |
| 67 | +/// use gpio_cdev::{AsyncLineEventHandle, Chip, EventRequestFlags, LineRequestFlags}; |
| 68 | +/// |
| 69 | +/// async fn print_events(line: u32) -> Result<()> { |
| 70 | +/// let mut chip = Chip::new("/dev/gpiochip0")?; |
| 71 | +/// let line = chip.get_line(line)?; |
| 72 | +/// let mut events = AsyncLineEventHandle::new(line.events( |
| 73 | +/// LineRequestFlags::INPUT, |
| 74 | +/// EventRequestFlags::BOTH_EDGES, |
| 75 | +/// "gpioevents", |
| 76 | +/// )?)?; |
| 77 | +/// |
| 78 | +/// loop { |
| 79 | +/// match events.next().await { |
| 80 | +/// Some(event) => println!("{:?}", event?), |
| 81 | +/// None => break, |
| 82 | +/// }; |
| 83 | +/// } |
| 84 | +/// |
| 85 | +/// Ok(()) |
| 86 | +/// } |
| 87 | +/// |
| 88 | +/// # #[tokio::main] |
| 89 | +/// # async fn main() { |
| 90 | +/// # print_events(42).await.unwrap(); |
| 91 | +/// # } |
| 92 | +/// ``` |
| 93 | +pub struct AsyncLineEventHandle { |
| 94 | + evented: PollEvented<PollWrapper>, |
| 95 | +} |
| 96 | + |
| 97 | +impl AsyncLineEventHandle { |
| 98 | + /// Wraps the specified `LineEventHandle`. |
| 99 | + /// |
| 100 | + /// # Arguments |
| 101 | + /// |
| 102 | + /// * `handle` - handle to be wrapped. |
| 103 | + pub fn new(handle: LineEventHandle) -> Result<AsyncLineEventHandle> { |
| 104 | + // The file descriptor needs to be configured for non-blocking I/O for PollEvented to work. |
| 105 | + let fd = handle.file.as_raw_fd(); |
| 106 | + unsafe { |
| 107 | + let flags = libc::fcntl(fd, libc::F_GETFL, 0); |
| 108 | + libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK); |
| 109 | + } |
| 110 | + |
| 111 | + Ok(AsyncLineEventHandle { |
| 112 | + evented: PollEvented::new(PollWrapper { handle })?, |
| 113 | + }) |
| 114 | + } |
| 115 | +} |
| 116 | + |
| 117 | +impl Stream for AsyncLineEventHandle { |
| 118 | + type Item = Result<LineEvent>; |
| 119 | + |
| 120 | + fn poll_next(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Option<Self::Item>> { |
| 121 | + let ready = Ready::readable(); |
| 122 | + if let Err(e) = ready!(self.evented.poll_read_ready(cx, ready)) { |
| 123 | + return Poll::Ready(Some(Err(e.into()))); |
| 124 | + } |
| 125 | + |
| 126 | + // TODO: This code should not be duplicated here. |
| 127 | + let mut data: ffi::gpioevent_data = unsafe { mem::zeroed() }; |
| 128 | + let mut data_as_buf = unsafe { |
| 129 | + slice::from_raw_parts_mut( |
| 130 | + &mut data as *mut ffi::gpioevent_data as *mut u8, |
| 131 | + mem::size_of::<ffi::gpioevent_data>(), |
| 132 | + ) |
| 133 | + }; |
| 134 | + match nix::unistd::read( |
| 135 | + self.evented.get_ref().handle.file.as_raw_fd(), |
| 136 | + &mut data_as_buf, |
| 137 | + ) { |
| 138 | + Ok(bytes_read) => { |
| 139 | + if bytes_read != mem::size_of::<ffi::gpioevent_data>() { |
| 140 | + let e = nix::Error::Sys(nix::errno::Errno::EIO); |
| 141 | + Poll::Ready(Some(Err(event_err(e)))) |
| 142 | + } else { |
| 143 | + Poll::Ready(Some(Ok(LineEvent(data)))) |
| 144 | + } |
| 145 | + } |
| 146 | + Err(nix::Error::Sys(nix::errno::Errno::EAGAIN)) => { |
| 147 | + self.evented.clear_read_ready(cx, ready)?; |
| 148 | + Poll::Pending |
| 149 | + } |
| 150 | + Err(e) => Poll::Ready(Some(Err(event_err(e)))), |
| 151 | + } |
| 152 | + } |
| 153 | +} |
| 154 | + |
| 155 | +impl AsRef<LineEventHandle> for AsyncLineEventHandle { |
| 156 | + fn as_ref(&self) -> &LineEventHandle { |
| 157 | + &self.evented.get_ref().handle |
| 158 | + } |
| 159 | +} |
0 commit comments