rust: bssl-tls: std and tokio transport integration This adds direct battery-included integration with `std` and `tokio`. Bug: 479599893 Signed-off-by: Xiangfei Ding <xfding@google.com> Change-Id: Id9ce290478ebbaddc80243db983025d46a6a6964 Reviewed-on: https://boringssl-review.googlesource.com/c/boringssl/+/92189 Reviewed-by: Rudolf Polzer <rpolzer@google.com> Presubmit-BoringSSL-Verified: boringssl-scoped@luci-project-accounts.iam.gserviceaccount.com <boringssl-scoped@luci-project-accounts.iam.gserviceaccount.com> Reviewed-by: Adam Langley <agl@google.com>
diff --git a/rust/bssl-tls/src/io.rs b/rust/bssl-tls/src/io.rs index 81ec098..ce0c5d7 100644 --- a/rust/bssl-tls/src/io.rs +++ b/rust/bssl-tls/src/io.rs
@@ -47,6 +47,13 @@ #[cfg(feature = "std")] pub mod stdio; +/// Synchronous I/O adapters. +pub mod sync_io; +/// Tokio-based async I/O adapters. +#[cfg(feature = "tokio_net")] +pub mod tokio; +#[cfg(all(unix, feature = "std"))] +pub mod unix; /// A wrapper around a `dyn AbstractSocket`, delegating BIO methods to the /// underlying `AbstractSocket` implementations.
diff --git a/rust/bssl-tls/src/io/sync_io.rs b/rust/bssl-tls/src/io/sync_io.rs new file mode 100644 index 0000000..83c1b6a --- /dev/null +++ b/rust/bssl-tls/src/io/sync_io.rs
@@ -0,0 +1,159 @@ +// Copyright 2026 The BoringSSL Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::io::stdio::PollFor; +use crate::io::{AbstractReader, AbstractSocket, AbstractSocketResult, AbstractWriter}; +use std::{ + io, + marker::PhantomData, + task::{Context, Poll}, +}; + +#[macro_export] +#[doc(hidden)] +macro_rules! retry_on_interrupt { + ($e:ident) => { + if matches!($e.kind(), ::std::io::ErrorKind::Interrupted) { + continue; + } else { + $crate::io::sync_io::translate_stdio_err($e) + } + }; +} + +/// `std::io::Read` and `std::io::Write` wrapper. +pub struct StdIoWithReactor<Io, Reactor> { + io: Io, + reactor: Reactor, + _p: PhantomData<fn() -> Reactor>, +} + +impl<Io, Reactor> StdIoWithReactor<Io, Reactor> { + /// Construct a new I/O wrapper + pub fn new(io: Io, reactor: Reactor) -> Self { + Self { + io, + reactor, + _p: PhantomData, + } + } +} + +impl<Io: io::Read + Send, R: PollFor<Io> + Send> AbstractReader for StdIoWithReactor<Io, R> { + fn read( + &mut self, + mut async_ctx: Option<&mut Context<'_>>, + buffer: &mut [u8], + ) -> AbstractSocketResult { + loop { + let res = match <Io as io::Read>::read(&mut self.io, buffer) { + Ok(0) if !buffer.is_empty() => return AbstractSocketResult::EndOfStream, + Ok(bytes) => return AbstractSocketResult::Ok(bytes), + Err(e) => retry_on_interrupt!(e), + }; + if let Some(async_ctx) = &mut async_ctx + && matches!(res, AbstractSocketResult::Retry) + { + match self.reactor.poll_read(async_ctx) { + Poll::Ready(Err(e)) => return AbstractSocketResult::Err(Box::new(e)), + Poll::Pending => return AbstractSocketResult::Retry, + Poll::Ready(Ok(())) => continue, + } + } + return res; + } + } +} + +impl<Io: io::Write + Send, R: PollFor<Io> + Send> AbstractWriter for StdIoWithReactor<Io, R> { + fn write( + &mut self, + mut async_ctx: Option<&mut Context<'_>>, + buffer: &[u8], + ) -> AbstractSocketResult { + loop { + let res = match <Io as io::Write>::write(&mut self.io, buffer) { + Ok(bytes) => return AbstractSocketResult::Ok(bytes), + Err(e) => retry_on_interrupt!(e), + }; + if let Some(async_ctx) = &mut async_ctx + && matches!(res, AbstractSocketResult::Retry) + { + match self.reactor.poll_write(async_ctx) { + Poll::Ready(Err(e)) => return AbstractSocketResult::Err(Box::new(e)), + Poll::Pending => return AbstractSocketResult::Retry, + Poll::Ready(Ok(())) => continue, + } + } + return res; + } + } + + fn flush(&mut self, mut async_ctx: Option<&mut Context<'_>>) -> AbstractSocketResult { + loop { + let res = match <Io as io::Write>::flush(&mut self.io) { + Ok(_) => return AbstractSocketResult::Ok(0), + Err(e) => retry_on_interrupt!(e), + }; + if let Some(async_ctx) = &mut async_ctx + && matches!(res, AbstractSocketResult::Retry) + { + match self.reactor.poll_write(async_ctx) { + Poll::Ready(Err(e)) => return AbstractSocketResult::Err(Box::new(e)), + Poll::Pending => return AbstractSocketResult::Retry, + Poll::Ready(Ok(())) => continue, + } + } + return res; + } + } +} + +impl<Io: io::Read + io::Write + Send, R: PollFor<Io> + Send> AbstractSocket + for StdIoWithReactor<Io, R> +{ +} + +/// Translates a `std::io::Error` into an `AbstractSocketResult`. +pub(crate) fn translate_stdio_err(err: io::Error) -> AbstractSocketResult { + match err.kind() { + io::ErrorKind::WouldBlock => AbstractSocketResult::Retry, + io::ErrorKind::ConnectionReset + | io::ErrorKind::ConnectionRefused + | io::ErrorKind::ConnectionAborted + | io::ErrorKind::BrokenPipe + | io::ErrorKind::NotConnected + | io::ErrorKind::UnexpectedEof => AbstractSocketResult::EndOfStream, + _ => AbstractSocketResult::Err(Box::new(err)), + } +} + +/// Reactor is absent. +/// Polling returns immediately with success, turning any I/O operation into a blocking call. +/// +/// # Warning ⚠️ +/// +/// If the underlying I/O is non-blocking, this reactor will lead to busy-waiting and excessive +/// CPU usage. +pub struct NoAsync; + +impl<Io> PollFor<Io> for NoAsync { + fn poll_read(&mut self, _: &mut Context<'_>) -> Poll<Result<(), io::Error>> { + Poll::Ready(Ok(())) + } + + fn poll_write(&mut self, _: &mut Context<'_>) -> Poll<Result<(), io::Error>> { + Poll::Ready(Ok(())) + } +}
diff --git a/rust/bssl-tls/src/io/tokio.rs b/rust/bssl-tls/src/io/tokio.rs new file mode 100644 index 0000000..4953edf --- /dev/null +++ b/rust/bssl-tls/src/io/tokio.rs
@@ -0,0 +1,131 @@ +// Copyright 2026 The BoringSSL Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#![cfg(feature = "tokio_net")] + +use std::{ + pin::Pin, + task::{ + Context, + Poll, // + }, // +}; + +use tokio::io::{ + AsyncRead, + AsyncWrite, + ReadBuf, // +}; + +use crate::io::{ + AbstractReader, + AbstractSocket, + AbstractSocketResult, + AbstractWriter, + NoAsyncContext, // +}; + +/// IO object implementing [`tokio::io::AsyncRead`] or [`tokio::io::AsyncWrite`] protocol. +pub struct TokioIo<T>(pub T); + +fn tokio_async_read<T: AsyncRead>( + mut this: Pin<&mut T>, + ctx: &mut Context<'_>, + buffer: &mut [u8], +) -> AbstractSocketResult { + let mut buf = ReadBuf::new(buffer); + loop { + return match this.as_mut().poll_read(ctx, &mut buf) { + Poll::Ready(Ok(())) => { + if buf.filled().is_empty() && buf.remaining() > 0 { + AbstractSocketResult::EndOfStream + } else { + AbstractSocketResult::Ok(buf.filled().len()) + } + } + Poll::Pending => AbstractSocketResult::Retry, + Poll::Ready(Err(e)) => crate::retry_on_interrupt!(e), + }; + } +} + +fn tokio_async_write<T: AsyncWrite>( + mut this: Pin<&mut T>, + ctx: &mut Context<'_>, + buffer: &[u8], +) -> AbstractSocketResult { + loop { + return match this.as_mut().poll_write(ctx, buffer) { + Poll::Ready(Ok(bytes)) => { + if buffer.is_empty() { + AbstractSocketResult::Ok(0) + } else if bytes == 0 { + AbstractSocketResult::EndOfStream + } else { + AbstractSocketResult::Ok(bytes) + } + } + Poll::Pending => AbstractSocketResult::Retry, + Poll::Ready(Err(e)) => crate::retry_on_interrupt!(e), + }; + } +} + +fn tokio_async_flush<T: AsyncWrite>( + mut this: Pin<&mut T>, + ctx: &mut Context<'_>, +) -> AbstractSocketResult { + loop { + return match this.as_mut().poll_flush(ctx) { + Poll::Ready(Ok(())) => AbstractSocketResult::Ok(0), + Poll::Pending => AbstractSocketResult::Retry, + Poll::Ready(Err(e)) => crate::retry_on_interrupt!(e), + }; + } +} + +impl<T: AsyncRead + Send + Unpin> AbstractReader for TokioIo<T> { + fn read( + &mut self, + async_ctx: Option<&mut Context<'_>>, + buffer: &mut [u8], + ) -> AbstractSocketResult { + let Some(ctx) = async_ctx else { + return AbstractSocketResult::Err(Box::new(NoAsyncContext)); + }; + tokio_async_read(Pin::new(&mut self.0), ctx, buffer) + } +} + +impl<T: AsyncWrite + Send + Unpin> AbstractWriter for TokioIo<T> { + fn write( + &mut self, + async_ctx: Option<&mut Context<'_>>, + buffer: &[u8], + ) -> AbstractSocketResult { + let Some(ctx) = async_ctx else { + return AbstractSocketResult::Err(Box::new(NoAsyncContext)); + }; + tokio_async_write(Pin::new(&mut self.0), ctx, buffer) + } + + fn flush(&mut self, async_ctx: Option<&mut Context<'_>>) -> AbstractSocketResult { + let Some(ctx) = async_ctx else { + return AbstractSocketResult::Err(Box::new(NoAsyncContext)); + }; + tokio_async_flush(Pin::new(&mut self.0), ctx) + } +} + +impl<T: AsyncRead + AsyncWrite + Send + Unpin> AbstractSocket for TokioIo<T> {}
diff --git a/rust/bssl-tls/src/io/unix.rs b/rust/bssl-tls/src/io/unix.rs new file mode 100644 index 0000000..1e0d886 --- /dev/null +++ b/rust/bssl-tls/src/io/unix.rs
@@ -0,0 +1,230 @@ +// Copyright 2026 The BoringSSL Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! I/O model for Unix-like systems +//! +//! It is strongly recommended that the non-blocking I/O is enabled by calling +//! [`std::net::TcpStream::set_nonblocking`] to `true`. + +use std::{ + io, + os::{ + fd::{AsRawFd, RawFd}, + unix::net::UnixDatagram, + }, + task::{Context, Poll}, +}; + +#[cfg(feature = "tokio_net")] +pub use crate::io::tokio::TokioIo; +#[cfg(feature = "tokio_net")] +pub use tokio_impl::TokioOverFd; + +use crate::ffi::{mut_slice_into_ffi_raw_parts, slice_into_ffi_raw_parts}; +use crate::io::stdio::{DatagramSocket, PollFor}; +use crate::io::sync_io::translate_stdio_err; + +use super::{AbstractReader, AbstractSocket, AbstractSocketResult, AbstractWriter}; + +#[cfg(feature = "tokio_net")] +mod tokio_impl; + +// ============ +// Datagrams +// ============ + +/// A datagram socket. +pub struct StdDatagram<Socket, Reactor> { + reactor: Reactor, + socket: Socket, +} + +impl<S, R> StdDatagram<S, R> { + /// Construct a new Unix Datagram. + pub fn new(socket: S, reactor: R) -> Self { + Self { socket, reactor } + } +} + +impl<Socket: DatagramSocket, Reactor: PollFor<Socket> + Send> AbstractReader + for StdDatagram<Socket, Reactor> +{ + fn read( + &mut self, + mut async_ctx: Option<&mut Context<'_>>, + buffer: &mut [u8], + ) -> AbstractSocketResult { + loop { + match (self.socket.recv(buffer), async_ctx.as_mut()) { + (AbstractSocketResult::Retry, Some(ctx)) => match self.reactor.poll_read(ctx) { + Poll::Pending => return AbstractSocketResult::Retry, + Poll::Ready(Ok(_)) => {} + Poll::Ready(Err(e)) => { + return AbstractSocketResult::Err(Box::new(e)); + } + }, + (res, _) => return res, + } + } + } +} + +impl<Socket: DatagramSocket, Reactor: PollFor<Socket> + Send> AbstractWriter + for StdDatagram<Socket, Reactor> +{ + fn write( + &mut self, + mut async_ctx: Option<&mut Context<'_>>, + buffer: &[u8], + ) -> AbstractSocketResult { + loop { + match (self.socket.send(buffer), async_ctx.as_mut()) { + (AbstractSocketResult::Retry, Some(ctx)) => { + if matches!(self.reactor.poll_write(ctx), Poll::Pending) { + return AbstractSocketResult::Retry; + } + } + (res, _) => return res, + } + } + } + + fn flush(&mut self, _: Option<&mut Context<'_>>) -> AbstractSocketResult { + AbstractSocketResult::Ok(0) + } +} + +impl<Socket: DatagramSocket, Reactor: PollFor<Socket> + Send> AbstractSocket + for StdDatagram<Socket, Reactor> +{ +} + +impl DatagramSocket for UnixDatagram { + fn send(&mut self, datagram: &[u8]) -> AbstractSocketResult { + loop { + return match UnixDatagram::send(self, datagram) { + Ok(bytes) => AbstractSocketResult::Ok(bytes), + Err(e) => crate::retry_on_interrupt!(e), + }; + } + } + + fn recv(&mut self, datagram: &mut [u8]) -> AbstractSocketResult { + loop { + return match UnixDatagram::recv(self, datagram) { + Ok(bytes) => AbstractSocketResult::Ok(bytes), + Err(e) if matches!(e.kind(), io::ErrorKind::Interrupted) => continue, + Err(e) => crate::retry_on_interrupt!(e), + }; + } + } +} + +/// A wrapper to operate the IO object through file descriptor. +/// +/// # Notice to **BSD** datagram socket users +/// +/// Whether a `SIGPIPE` signal will be raised is controlled at socket configuration time with +/// `setsocketopt` plus `SO_NOSIGPIPE`. +/// If the signal should be suppressed, it should be configured before using this type. +pub struct UseFd<Io: AsRawFd>(pub Io); + +#[cfg(feature = "libc")] +impl<Io: AsRawFd> io::Read for UseFd<Io> { + fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> { + let (ptr, len) = mut_slice_into_ffi_raw_parts(buf); + let ret = unsafe { + // Safety: `ptr` has been valid with the sufficient capacity of length `len`. + libc::read(self.as_raw_fd(), ptr as _, len) + }; + if ret < 0 { + Err(io::Error::last_os_error()) + } else { + Ok(ret as usize) + } + } +} + +#[cfg(feature = "libc")] +impl<Io: AsRawFd> io::Write for UseFd<Io> { + fn write(&mut self, buf: &[u8]) -> io::Result<usize> { + let (ptr, len) = slice_into_ffi_raw_parts(buf); + let ret = unsafe { + // Safety: `ptr` has been valid for length `len`. + libc::write(self.as_raw_fd(), ptr as _, len) + }; + if ret < 0 { + Err(io::Error::last_os_error()) + } else { + Ok(ret as usize) + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +#[cfg(feature = "libc")] +impl<Io: AsRawFd + Send> DatagramSocket for UseFd<Io> { + fn send(&mut self, datagram: &[u8]) -> AbstractSocketResult { + let (buf, len) = slice_into_ffi_raw_parts(datagram); + #[cfg(any(target_os = "linux", target_os = "android"))] + let flag = libc::MSG_NOSIGNAL; + #[cfg(any( + target_os = "freebsd", + target_os = "netbsd", + target_os = "openbsd", + target_os = "macos", + target_os = "ios", + ))] + let flag = 0; + #[cfg(any(windows, target_os = "none"))] + let flag = 0; + loop { + let rc = unsafe { + // Safety: the socket file descriptor is exclusively owned. + libc::send(self.as_raw_fd(), buf as _, len, flag) + }; + return if rc < 0 { + let err = io::Error::last_os_error(); + crate::retry_on_interrupt!(err) + } else { + AbstractSocketResult::Ok(rc as usize) + }; + } + } + + fn recv(&mut self, datagram: &mut [u8]) -> AbstractSocketResult { + let (buf, len) = mut_slice_into_ffi_raw_parts(datagram); + loop { + let rc = unsafe { + // Safety: the socket file descriptor is exclusively owned. + libc::recv(self.as_raw_fd(), buf as _, len, 0) + }; + return if rc < 0 { + let err = io::Error::last_os_error(); + crate::retry_on_interrupt!(err) + } else { + AbstractSocketResult::Ok(rc as usize) + }; + } + } +} + +impl<Io: AsRawFd> AsRawFd for UseFd<Io> { + fn as_raw_fd(&self) -> RawFd { + self.0.as_raw_fd() + } +}
diff --git a/rust/bssl-tls/src/io/unix/tokio_impl.rs b/rust/bssl-tls/src/io/unix/tokio_impl.rs new file mode 100644 index 0000000..037d7af --- /dev/null +++ b/rust/bssl-tls/src/io/unix/tokio_impl.rs
@@ -0,0 +1,106 @@ +use std::{ + io, + os::fd::{AsRawFd, RawFd}, + task::{Context, Poll, ready}, +}; + +use tokio::io::{Interest, ReadBuf, Ready, unix::AsyncFd}; + +use crate::io::{ + AbstractReader, AbstractSocket, AbstractSocketResult, AbstractWriter, NoAsyncContext, + stdio::PollFor, + unix::{StdDatagram, UseFd, translate_stdio_err}, +}; + +/// Reactor that operate over file descriptor. +pub struct TokioOverFd(AsyncFd<RawFd>); + +impl<T: AsRawFd> StdDatagram<UseFd<T>, TokioOverFd> { + /// Construct a datagram IO object driven by [`tokio`]. + pub fn new_with_tokio(inner: T) -> Result<Self, io::Error> { + let reactor = TokioOverFd::new(inner.as_raw_fd())?; + let fd = UseFd(inner); + Ok(Self::new(fd, reactor)) + } +} + +impl TokioOverFd { + /// A trivial constructor to signal use of `tokio` reactor and register events with + /// file descriptors. + pub fn new(fd: RawFd) -> Result<Self, io::Error> { + Ok(Self(AsyncFd::try_with_interest( + fd, + Interest::READABLE | Interest::WRITABLE | Interest::ERROR, + )?)) + } +} + +impl<T> PollFor<T> for TokioOverFd { + fn poll_read(&mut self, async_ctx: &mut Context<'_>) -> Poll<Result<(), io::Error>> { + match ready!(self.0.poll_read_ready_mut(async_ctx)) { + Ok(mut guard) => { + guard.clear_ready_matching(Ready::READABLE); + Poll::Ready(Ok(())) + } + Err(e) => Poll::Ready(Err(e)), + } + } + + fn poll_write(&mut self, async_ctx: &mut Context<'_>) -> Poll<Result<(), io::Error>> { + match ready!(self.0.poll_write_ready_mut(async_ctx)) { + Ok(mut guard) => { + guard.clear_ready_matching(Ready::WRITABLE); + Poll::Ready(Ok(())) + } + Err(e) => Poll::Ready(Err(e)), + } + } +} + +macro_rules! gen_impl_datagram { + ($ty:ty) => { + impl AbstractReader for $ty { + fn read( + &mut self, + async_ctx: Option<&mut Context<'_>>, + buffer: &mut [u8], + ) -> AbstractSocketResult { + let Some(cx) = async_ctx else { + return AbstractSocketResult::Err(Box::new(NoAsyncContext)); + }; + let mut buf = ReadBuf::new(buffer); + match self.poll_recv(cx, &mut buf) { + Poll::Pending => AbstractSocketResult::Retry, + Poll::Ready(Ok(_)) => AbstractSocketResult::Ok(buf.filled().len()), + Poll::Ready(Err(e)) => translate_stdio_err(e), + } + } + } + + impl AbstractWriter for $ty { + fn write( + &mut self, + async_ctx: Option<&mut Context<'_>>, + buf: &[u8], + ) -> AbstractSocketResult { + let Some(cx) = async_ctx else { + return AbstractSocketResult::Err(Box::new(NoAsyncContext)); + }; + match self.poll_send(cx, buf) { + Poll::Pending => AbstractSocketResult::Retry, + Poll::Ready(Ok(len)) => AbstractSocketResult::Ok(len), + Poll::Ready(Err(e)) => translate_stdio_err(e), + } + } + + fn flush(&mut self, _: Option<&mut Context<'_>>) -> AbstractSocketResult { + AbstractSocketResult::Ok(0) + } + } + + impl AbstractSocket for $ty {} + }; +} + +gen_impl_datagram!(tokio::net::UdpSocket); +gen_impl_datagram!(tokio::net::UnixDatagram);