Структура Receiver
pub struct Receiver { /* private fields */ }
net.Читающий конец канала Unix.
Его можно создать из файла FIFO с помощью OpenOptions::open_receiver.
Примеры
Получение сообщений из именованного канала в цикле:
use tokio::net::unix::pipe;
use tokio::io::{self, AsyncReadExt};
const FIFO_NAME: &str = "path/to/a/fifo";
let mut rx = pipe::OpenOptions::new().open_receiver(FIFO_NAME)?;
loop {
let mut msg = vec![0; 256];
match rx.read_exact(&mut msg).await {
Ok(_) => {
/* handle the message */
}
Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => {
// Writing end has been closed, we should reopen the pipe.
rx = pipe::OpenOptions::new().open_receiver(FIFO_NAME)?;
}
Err(e) => return Err(e.into()),
}
}В Linux можно использовать Receiver в режиме чтения-записи для надёжного чтения из именованного канала. В отличие от Receiver, открытого в режиме только для чтения, чтение из канала в режиме чтения-записи не завершится ошибкой UnexpectedEof, когда записывающий конец будет закрыт. Таким образом, Receiver может асинхронно ожидать, пока следующий отправитель откроет канал.
Не следует использовать функции, ожидающие EOF, такие как read_to_end, с Receiver в режиме чтения-записи, так как они могут ждать бесконечно. Receiver в этом режиме также удерживает открытым записывающий конец, что препятствует получению EOF.
Чтобы установить режим чтения-записи, можно использовать OpenOptions::read_write. Обратите внимание, что использование режима чтения-записи с файлами FIFO не определено стандартом POSIX и гарантированно работает только в Linux.
use tokio::net::unix::pipe;
use tokio::io::AsyncReadExt;
const FIFO_NAME: &str = "path/to/a/fifo";
let mut rx = pipe::OpenOptions::new()
.read_write(true)
.open_receiver(FIFO_NAME)?;
loop {
let mut msg = vec![0; 256];
rx.read_exact(&mut msg).await?;
/* handle the message */
}
Реализации
impl Receiver
pub fn from_file(file: File) -> Result<Receiver>
Создаёт новый Receiver из File.
Эта функция предназначена для создания канала из File, представляющего специальный FIFO-файл. Она проверяет, является ли файл каналом и открыт ли он для чтения, устанавливает неблокирующий режим и выполняет преобразование.
Ошибки
Завершается ошибкой с io::ErrorKind::InvalidInput, если файл не является каналом или не открыт для чтения. Также функция завершается ошибкой при возникновении любой стандартной ошибки ОС.
Паники
Эта функция вызывает панику, если её вызывают вне среды выполнения с включённым вводом-выводом.
Обычно среда выполнения устанавливается неявно, когда эта функция вызывается из future, выполняемого средой выполнения tokio; в противном случае среду выполнения можно задать явно с помощью функции Runtime::enter.
pub fn from_owned_fd(owned_fd: OwnedFd) -> Result<Receiver>
Создаёт новый Receiver из OwnedFd.
Эта функция предназначена для создания канала из OwnedFd, представляющего анонимный канал или специальный FIFO-файл. Она проверяет, является ли файловый дескриптор каналом и открыт ли он для чтения, устанавливает неблокирующий режим и выполняет преобразование.
Ошибки
Завершается ошибкой с io::ErrorKind::InvalidInput, если файловый дескриптор не является каналом или не открыт для чтения. Также функция завершается ошибкой при возникновении любой стандартной ошибки ОС.
Паники
Эта функция вызывает панику, если её вызывают вне среды выполнения с включённым вводом-выводом.
Обычно среда выполнения устанавливается неявно, когда эта функция вызывается из future, выполняемого средой выполнения tokio; в противном случае среду выполнения можно задать явно с помощью функции Runtime::enter.
pub fn from_file_unchecked(file: File) -> Result<Receiver>
Создаёт новый Receiver из File без проверки свойств канала.
Эта функция предназначена для создания канала из File, представляющего специальный FIFO-файл. При преобразовании не проверяются свойства нижележащего файла; пользователь должен самостоятельно убедиться, что файл открыт для чтения, представляет собой канал и настроен в неблокирующем режиме.
Примеры
use tokio::net::unix::pipe;
use std::fs::OpenOptions;
use std::os::unix::fs::{FileTypeExt, OpenOptionsExt};
const FIFO_NAME: &str = "path/to/a/fifo";
let file = OpenOptions::new()
.read(true)
.custom_flags(libc::O_NONBLOCK)
.open(FIFO_NAME)?;
if file.metadata()?.file_type().is_fifo() {
let rx = pipe::Receiver::from_file_unchecked(file)?;
/* use the Receiver */
}Паники
Эта функция вызывает панику, если её вызывают вне среды выполнения с включённым вводом-выводом.
Обычно среда выполнения устанавливается неявно, когда эта функция вызывается из future, выполняемого средой выполнения tokio; в противном случае среду выполнения можно задать явно с помощью функции Runtime::enter.
pub fn from_owned_fd_unchecked(owned_fd: OwnedFd) -> Result<Receiver>
Создаёт новый Receiver из OwnedFd без проверки свойств канала.
Эта функция предназначена для создания канала из OwnedFd, представляющего анонимный канал или специальный FIFO-файл. При преобразовании не проверяются свойства нижележащего канала; пользователь должен самостоятельно убедиться, что файловый дескриптор представляет читающий конец канала и что канал настроен в неблокирующем режиме.
Паники
Эта функция вызывает панику, если её вызывают вне среды выполнения с включённым вводом-выводом.
Обычно среда выполнения устанавливается неявно, когда эта функция вызывается из future, выполняемого средой выполнения tokio; в противном случае среду выполнения можно задать явно с помощью функции Runtime::enter.
pub async fn ready(&self, interest: Interest) -> Result<Ready>
Ожидает наступления любого из запрошенных состояний готовности.
Эту функцию можно использовать вместо readable(), чтобы проверить возвращённый набор состояний готовности на наличие событий Ready::READABLE и Ready::READ_CLOSED.
Функция может завершиться, даже если канал не готов. Это ложное срабатывание, и попытка выполнить операцию вернёт io::ErrorKind::WouldBlock. Функция также может вернуть пустой набор Ready, поэтому всегда проверяйте возвращённое значение и, возможно, ожидайте снова, если запрошенные состояния не установлены.
Безопасность при отмене
Этот метод безопасен при отмене. После возникновения события готовности метод будет немедленно возвращать результат, пока это событие готовности не будет обработано попыткой чтения, завершившейся с ошибкой WouldBlock или Poll::Pending.
pub async fn readable(&self) -> Result<()>
Ожидает, пока канал не станет доступен для чтения.
Эта функция эквивалентна ready(Interest::READABLE) и обычно используется вместе с try_read().
Примеры
use tokio::net::unix::pipe;
use std::io;
#[tokio::main]
async fn main() -> io::Result<()> {
// Open a reading end of a fifo
let rx = pipe::OpenOptions::new().open_receiver("path/to/a/fifo")?;
let mut msg = vec![0; 1024];
loop {
// Wait for the pipe to be readable
rx.readable().await?;
// Try to read data, this may still fail with `WouldBlock`
// if the readiness event is a false positive.
match rx.try_read(&mut msg) {
Ok(n) => {
msg.truncate(n);
break;
}
Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
continue;
}
Err(e) => {
return Err(e.into());
}
}
}
println!("GOT = {:?}", msg);
Ok(())
}pub fn poll_read_ready(&self, cx: &mut Context<'_>) -> Poll<Result<()>>
Проверяет готовность к чтению.
Если канал в данный момент не готов к чтению, этот метод сохранит копию Waker из переданного Context. Когда канал будет готов к чтению, для пробуждающего объекта будет вызван Waker::wake.
Обратите внимание: при нескольких вызовах poll_read_ready или poll_read пробуждение будет запланировано только для Waker из Context, переданного при последнем вызове.
Эта функция предназначена для случаев, когда создать и закрепить future с помощью readable невозможно. Если есть такая возможность, предпочтительно использовать readable, поскольку это позволяет выполнять опрос одновременно из нескольких задач.
Возвращаемое значение
Функция возвращает:
-
Poll::Pending, если канал не готов к чтению. -
Poll::Ready(Ok(())), если канал готов к чтению. -
Poll::Ready(Err(e)), если произошла ошибка.
Ошибки
Эта функция может столкнуться с любой стандартной ошибкой ввода-вывода, кроме WouldBlock.
pub fn try_read(&self, buf: &mut [u8]) -> Result<usize>
Пытается прочитать данные из канала в предоставленный буфер и возвращает количество прочитанных байтов.
Читает все ожидающие данные из канала, но не ждёт поступления новых. В случае успеха возвращает количество прочитанных байтов. Поскольку try_read() работает в неблокирующем режиме, буфер не нужно сохранять в асинхронной задаче, и он может целиком находиться в стеке.
Эта функция обычно используется вместе с readable().
Возвращаемое значение
Если данные прочитаны успешно, возвращается Ok(n), где n — количество прочитанных байтов. Если n равно 0, это может указывать на один из двух сценариев:
- Записывающий конец канала закрыт и больше не будет записывать данные.
- Длина указанного буфера равна 0 байт.
Если канал не готов к чтению данных, возвращается Err(io::ErrorKind::WouldBlock).
Примечания
Чтобы избежать ненужных системных вызовов, чтение будет предприниматься только в том случае, если ОС сообщила Tokio, что канал стал доступен для чтения. Поэтому try_read() может завершиться ошибкой WouldBlock, если Tokio ещё не получил от ОС уведомление о том, что канал стал доступен для чтения.
Примеры
use tokio::net::unix::pipe;
use std::io;
#[tokio::main]
async fn main() -> io::Result<()> {
// Open a reading end of a fifo
let rx = pipe::OpenOptions::new().open_receiver("path/to/a/fifo")?;
let mut msg = vec![0; 1024];
loop {
// Wait for the pipe to be readable
rx.readable().await?;
// Try to read data, this may still fail with `WouldBlock`
// if the readiness event is a false positive.
match rx.try_read(&mut msg) {
Ok(n) => {
msg.truncate(n);
break;
}
Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
continue;
}
Err(e) => {
return Err(e.into());
}
}
}
println!("GOT = {:?}", msg);
Ok(())
}pub fn try_read_vectored(&self, bufs: &mut [IoSliceMut<'_>]) -> Result<usize>
Пытается прочитать данные из канала в предоставленные буферы и возвращает количество прочитанных байтов.
Данные копируются в буферы по порядку; последний буфер может быть заполнен лишь частично. Этот метод эквивалентен одному вызову try_read() с объединёнными буферами.
Считывает все ожидающие данные из канала, но не ждёт поступления новых данных. В случае успеха возвращает количество прочитанных байтов. Поскольку try_read_vectored() работает в неблокирующем режиме, буфер не требуется хранить в асинхронной задаче — он может целиком находиться в стеке.
Обычно эту функцию используют вместе с readable().
Возвращаемое значение
Если данные успешно прочитаны, возвращается Ok(n), где n — количество прочитанных байтов. Ok(0) означает, что записывающий конец канала закрыт и больше не будет записывать данные. Если канал не готов к чтению данных, возвращается Err(io::ErrorKind::WouldBlock).
Примечания
Чтобы избежать ненужных системных вызовов, операция чтения выполняется только в том случае, если ОС сообщила Tokio, что канал стал доступен для чтения. Поэтому try_read_vectored() может завершиться ошибкой WouldBlock, если Tokio ещё не получил от ОС уведомление о том, что канал стал доступен для чтения.
Примеры
use tokio::net::unix::pipe;
use std::io;
#[tokio::main]
async fn main() -> io::Result<()> {
// Open a reading end of a fifo
let rx = pipe::OpenOptions::new().open_receiver("path/to/a/fifo")?;
loop {
// Wait for the pipe to be readable
rx.readable().await?;
// Creating the buffer **after** the `await` prevents it from
// being stored in the async task.
let mut buf_a = [0; 512];
let mut buf_b = [0; 1024];
let mut bufs = [
io::IoSliceMut::new(&mut buf_a),
io::IoSliceMut::new(&mut buf_b),
];
// Try to read data, this may still fail with `WouldBlock`
// if the readiness event is a false positive.
match rx.try_read_vectored(&mut bufs) {
Ok(0) => break,
Ok(n) => {
println!("read {} bytes", n);
}
Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
continue;
}
Err(e) => {
return Err(e.into());
}
}
}
Ok(())
}pub fn try_io<R>(&self, f: impl FnOnce() -> Result<R>) -> Result<R>
Пытается выполнить чтение из сокета с помощью предоставленной пользователем операции ввода-вывода.
Если сокет готов, вызывается предоставленная замыкающая функция. Она должна попытаться выполнить операцию ввода-вывода над сокетом, вручную вызвав соответствующий системный вызов. Если операция завершается ошибкой, поскольку сокет фактически не готов, замыкающая функция должна вернуть ошибку WouldBlock, после чего флаг готовности сбрасывается. Затем try_io возвращает значение, возвращённое замыкающей функцией.
Если сокет не готов, замыкающая функция не вызывается, а возвращается ошибка WouldBlock.
Замыкающая функция должна возвращать ошибку WouldBlock только в том случае, если она выполнила операцию ввода-вывода над сокетом, которая завершилась неудачей из-за неготовности сокета. Возврат ошибки WouldBlock в любой другой ситуации приведёт к ошибочному сбросу флага готовности, из-за чего сокет может работать неправильно.
Замыкающая функция не должна выполнять операцию ввода-вывода с помощью методов, определённых для типа Tokio pipe::Receiver, так как это нарушит работу флага готовности и может привести к неправильной работе сокета.
Обычно эту функцию используют вместе с readable() или ready().
pub fn try_read_buf<B: BufMut>(&self, buf: &mut B) -> Result<usize>
io-util.Пытается прочитать данные из канала в предоставленный буфер, сдвигая его внутренний курсор, и возвращает количество прочитанных байтов.
Считывает все ожидающие данные из канала, но не ждёт поступления новых данных. В случае успеха возвращает количество прочитанных байтов. Поскольку try_read_buf() работает в неблокирующем режиме, буфер не требуется хранить в асинхронной задаче — он может целиком находиться в стеке.
Обычно эту функцию используют вместе с readable() или ready().
Возвращаемое значение
Если данные успешно прочитаны, возвращается Ok(n), где n — количество прочитанных байтов. Ok(0) означает, что записывающий конец канала закрыт и больше не будет записывать данные. Если канал не готов к чтению данных, возвращается Err(io::ErrorKind::WouldBlock).
Примечания
Чтобы избежать ненужных системных вызовов, операция чтения выполняется только в том случае, если ОС сообщила Tokio, что канал стал доступен для чтения. Поэтому try_read_buf() может завершиться ошибкой WouldBlock, если Tokio ещё не получил от ОС уведомление о том, что канал стал доступен для чтения.
Примеры
use tokio::net::unix::pipe;
use std::io;
#[tokio::main]
async fn main() -> io::Result<()> {
// Open a reading end of a fifo
let rx = pipe::OpenOptions::new().open_receiver("path/to/a/fifo")?;
loop {
// Wait for the pipe to be readable
rx.readable().await?;
let mut buf = Vec::with_capacity(4096);
// Try to read data, this may still fail with `WouldBlock`
// if the readiness event is a false positive.
match rx.try_read_buf(&mut buf) {
Ok(0) => break,
Ok(n) => {
println!("read {} bytes", n);
}
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
continue;
}
Err(e) => {
return Err(e.into());
}
}
}
Ok(())
}pub fn into_blocking_fd(self) -> Result<OwnedFd>
Преобразует канал в OwnedFd в блокирующем режиме.
Эта функция удалит этот конец канала из цикла обработки событий, установит для него блокирующий режим и выполнит преобразование.
pub fn into_nonblocking_fd(self) -> Result<OwnedFd>
Преобразует канал в OwnedFd в неблокирующем режиме.
Эта функция удалит этот конец канала из цикла обработки событий и выполнит преобразование. Возвращаемый файловый дескриптор будет работать в неблокирующем режиме.
Реализации трейтов
impl AsFd for Receiver
fn as_fd(&self) -> BorrowedFd<'_>
Автоматические реализации трейтов
impl !Freeze for Receiver
impl RefUnwindSafe for Receiver
impl Send for Receiver
impl Sync for Receiver
impl Unpin for Receiver
impl UnsafeUnpin for Receiver
impl UnwindSafe for Receiver
Общие реализации
impl<R> AsyncReadExt for R
fn chain<R>(self, next: R) -> Chain<Self, R>
io-util.fn read<'a>(&'a mut self, buf: &'a mut [u8]) -> Read<'a, Self>where Self: Unpin,
io-util.fn read_buf<'a, B>(&'a mut self, buf: &'a mut B) -> ReadBuf<'a, Self, B>
io-util.fn read_exact<'a>(&'a mut self, buf: &'a mut [u8]) -> ReadExact<'a, Self>where Self: Unpin,
io-util.buf. Подробнее
fn read_u8(&mut self) -> ReadU8<&mut Self>where Self: Unpin,
io-util.fn read_i8(&mut self) -> ReadI8<&mut Self>where Self: Unpin,
io-util.fn read_u32(&mut self) -> ReadU32<&mut Self>where Self: Unpin,
io-util.fn read_i32(&mut self) -> ReadI32<&mut Self>where Self: Unpin,
io-util.fn read_u64(&mut self) -> ReadU64<&mut Self>where Self: Unpin,
io-util.fn read_i64(&mut self) -> ReadI64<&mut Self>where Self: Unpin,
io-util.fn read_u128(&mut self) -> ReadU128<&mut Self>where Self: Unpin,
io-util.fn read_i128(&mut self) -> ReadI128<&mut Self>where Self: Unpin,
io-util.fn read_f32(&mut self) -> ReadF32<&mut Self>where Self: Unpin,
io-util.fn read_f64(&mut self) -> ReadF64<&mut Self>where Self: Unpin,
io-util.fn read_u16_le(&mut self) -> ReadU16Le<&mut Self>where Self: Unpin,
io-util.fn read_i16_le(&mut self) -> ReadI16Le<&mut Self>where Self: Unpin,
io-util.fn read_u32_le(&mut self) -> ReadU32Le<&mut Self>where Self: Unpin,
io-util.fn read_i32_le(&mut self) -> ReadI32Le<&mut Self>where Self: Unpin,
io-util.fn read_u64_le(&mut self) -> ReadU64Le<&mut Self>where Self: Unpin,
io-util.fn read_i64_le(&mut self) -> ReadI64Le<&mut Self>where Self: Unpin,
io-util.fn read_u128_le(&mut self) -> ReadU128Le<&mut Self>where Self: Unpin,
io-util.fn read_i128_le(&mut self) -> ReadI128Le<&mut Self>where Self: Unpin,
io-util.fn read_f32_le(&mut self) -> ReadF32Le<&mut Self>where Self: Unpin,
io-util.fn read_f64_le(&mut self) -> ReadF64Le<&mut Self>where Self: Unpin,
io-util.fn read_to_end<'a>(&'a mut self, buf: &'a mut Vec<u8>) -> ReadToEnd<'a, Self>where Self: Unpin,
io-util.buf. Подробнее
fn read_to_string<'a>( &'a mut self, dst: &'a mut String, ) -> ReadToString<'a, Self>where Self: Unpin,
io-util.buf. Подробнее
impl<T> BorrowMut<T> for Twhere T: ?Sized,
fn borrow_mut(&mut self) -> &mut T
impl<T> Instrument for T
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
impl<T> WithSubscriber for T
fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ
Subscriber к этому типу, возвращая обёртку WithDispatch. Подробнее
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/net/unix/pipe/struct.Receiver.html