mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-03 00:00:05 +02:00
util: add lifetime parameter to ReusableBoxFuture (#3762)
Co-authored-by: Toby Lawrence <[email protected]>
This commit is contained in:
@@ -14,7 +14,7 @@ use std::task::{Context, Poll};
|
|||||||
/// [`Stream`]: trait@crate::Stream
|
/// [`Stream`]: trait@crate::Stream
|
||||||
#[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
|
#[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
|
||||||
pub struct BroadcastStream<T> {
|
pub struct BroadcastStream<T> {
|
||||||
inner: ReusableBoxFuture<(Result<T, RecvError>, Receiver<T>)>,
|
inner: ReusableBoxFuture<'static, (Result<T, RecvError>, Receiver<T>)>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// An error returned from the inner stream of a [`BroadcastStream`].
|
/// An error returned from the inner stream of a [`BroadcastStream`].
|
||||||
|
|||||||
@@ -49,7 +49,7 @@ use tokio::sync::watch::error::RecvError;
|
|||||||
/// [`Stream`]: trait@crate::Stream
|
/// [`Stream`]: trait@crate::Stream
|
||||||
#[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
|
#[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
|
||||||
pub struct WatchStream<T> {
|
pub struct WatchStream<T> {
|
||||||
inner: ReusableBoxFuture<(Result<(), RecvError>, Receiver<T>)>,
|
inner: ReusableBoxFuture<'static, (Result<(), RecvError>, Receiver<T>)>,
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn make_future<T: Clone + Send + Sync>(
|
async fn make_future<T: Clone + Send + Sync>(
|
||||||
|
|||||||
@@ -44,7 +44,7 @@ enum State<T> {
|
|||||||
pub struct PollSender<T> {
|
pub struct PollSender<T> {
|
||||||
sender: Option<Sender<T>>,
|
sender: Option<Sender<T>>,
|
||||||
state: State<T>,
|
state: State<T>,
|
||||||
acquire: ReusableBoxFuture<Result<OwnedPermit<T>, PollSendError<T>>>,
|
acquire: ReusableBoxFuture<'static, Result<OwnedPermit<T>, PollSendError<T>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
// Creates a future for acquiring a permit from the underlying channel. This is used to ensure
|
// Creates a future for acquiring a permit from the underlying channel. This is used to ensure
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ use super::ReusableBoxFuture;
|
|||||||
/// [`Semaphore`]: tokio::sync::Semaphore
|
/// [`Semaphore`]: tokio::sync::Semaphore
|
||||||
pub struct PollSemaphore {
|
pub struct PollSemaphore {
|
||||||
semaphore: Arc<Semaphore>,
|
semaphore: Arc<Semaphore>,
|
||||||
permit_fut: Option<ReusableBoxFuture<Result<OwnedSemaphorePermit, AcquireError>>>,
|
permit_fut: Option<ReusableBoxFuture<'static, Result<OwnedSemaphorePermit, AcquireError>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl PollSemaphore {
|
impl PollSemaphore {
|
||||||
|
|||||||
@@ -6,26 +6,23 @@ use std::ptr::{self, NonNull};
|
|||||||
use std::task::{Context, Poll};
|
use std::task::{Context, Poll};
|
||||||
use std::{fmt, panic};
|
use std::{fmt, panic};
|
||||||
|
|
||||||
/// A reusable `Pin<Box<dyn Future<Output = T> + Send>>`.
|
/// A reusable `Pin<Box<dyn Future<Output = T> + Send + 'a>>`.
|
||||||
///
|
///
|
||||||
/// This type lets you replace the future stored in the box without
|
/// This type lets you replace the future stored in the box without
|
||||||
/// reallocating when the size and alignment permits this.
|
/// reallocating when the size and alignment permits this.
|
||||||
pub struct ReusableBoxFuture<T> {
|
pub struct ReusableBoxFuture<'a, T> {
|
||||||
boxed: NonNull<dyn Future<Output = T> + Send>,
|
boxed: NonNull<dyn Future<Output = T> + Send + 'a>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T> ReusableBoxFuture<T> {
|
impl<'a, T> ReusableBoxFuture<'a, T> {
|
||||||
/// Create a new `ReusableBoxFuture<T>` containing the provided future.
|
/// Create a new `ReusableBoxFuture<T>` containing the provided future.
|
||||||
pub fn new<F>(future: F) -> Self
|
pub fn new<F>(future: F) -> Self
|
||||||
where
|
where
|
||||||
F: Future<Output = T> + Send + 'static,
|
F: Future<Output = T> + Send + 'a,
|
||||||
{
|
{
|
||||||
let boxed: Box<dyn Future<Output = T> + Send> = Box::new(future);
|
let boxed: Box<dyn Future<Output = T> + Send + 'a> = Box::new(future);
|
||||||
|
|
||||||
let boxed = Box::into_raw(boxed);
|
let boxed = NonNull::from(Box::leak(boxed));
|
||||||
|
|
||||||
// SAFETY: Box::into_raw does not return null pointers.
|
|
||||||
let boxed = unsafe { NonNull::new_unchecked(boxed) };
|
|
||||||
|
|
||||||
Self { boxed }
|
Self { boxed }
|
||||||
}
|
}
|
||||||
@@ -36,7 +33,7 @@ impl<T> ReusableBoxFuture<T> {
|
|||||||
/// different from the layout of the currently stored future.
|
/// different from the layout of the currently stored future.
|
||||||
pub fn set<F>(&mut self, future: F)
|
pub fn set<F>(&mut self, future: F)
|
||||||
where
|
where
|
||||||
F: Future<Output = T> + Send + 'static,
|
F: Future<Output = T> + Send + 'a,
|
||||||
{
|
{
|
||||||
if let Err(future) = self.try_set(future) {
|
if let Err(future) = self.try_set(future) {
|
||||||
*self = Self::new(future);
|
*self = Self::new(future);
|
||||||
@@ -50,7 +47,7 @@ impl<T> ReusableBoxFuture<T> {
|
|||||||
/// future.
|
/// future.
|
||||||
pub fn try_set<F>(&mut self, future: F) -> Result<(), F>
|
pub fn try_set<F>(&mut self, future: F) -> Result<(), F>
|
||||||
where
|
where
|
||||||
F: Future<Output = T> + Send + 'static,
|
F: Future<Output = T> + Send + 'a,
|
||||||
{
|
{
|
||||||
// SAFETY: The pointer is not dangling.
|
// SAFETY: The pointer is not dangling.
|
||||||
let self_layout = {
|
let self_layout = {
|
||||||
@@ -78,7 +75,7 @@ impl<T> ReusableBoxFuture<T> {
|
|||||||
/// same as `self.layout`.
|
/// same as `self.layout`.
|
||||||
unsafe fn set_same_layout<F>(&mut self, future: F)
|
unsafe fn set_same_layout<F>(&mut self, future: F)
|
||||||
where
|
where
|
||||||
F: Future<Output = T> + Send + 'static,
|
F: Future<Output = T> + Send + 'a,
|
||||||
{
|
{
|
||||||
// Drop the existing future, catching any panics.
|
// Drop the existing future, catching any panics.
|
||||||
let result = panic::catch_unwind(AssertUnwindSafe(|| {
|
let result = panic::catch_unwind(AssertUnwindSafe(|| {
|
||||||
@@ -116,7 +113,7 @@ impl<T> ReusableBoxFuture<T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T> Future for ReusableBoxFuture<T> {
|
impl<T> Future for ReusableBoxFuture<'_, T> {
|
||||||
type Output = T;
|
type Output = T;
|
||||||
|
|
||||||
/// Poll the future stored inside this box.
|
/// Poll the future stored inside this box.
|
||||||
@@ -125,18 +122,18 @@ impl<T> Future for ReusableBoxFuture<T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// The future stored inside ReusableBoxFuture<T> must be Send.
|
// The future stored inside ReusableBoxFuture<'_, T> must be Send.
|
||||||
unsafe impl<T> Send for ReusableBoxFuture<T> {}
|
unsafe impl<T> Send for ReusableBoxFuture<'_, T> {}
|
||||||
|
|
||||||
// The only method called on self.boxed is poll, which takes &mut self, so this
|
// The only method called on self.boxed is poll, which takes &mut self, so this
|
||||||
// struct being Sync does not permit any invalid access to the Future, even if
|
// struct being Sync does not permit any invalid access to the Future, even if
|
||||||
// the future is not Sync.
|
// the future is not Sync.
|
||||||
unsafe impl<T> Sync for ReusableBoxFuture<T> {}
|
unsafe impl<T> Sync for ReusableBoxFuture<'_, T> {}
|
||||||
|
|
||||||
// Just like a Pin<Box<dyn Future>> is always Unpin, so is this type.
|
// Just like a Pin<Box<dyn Future>> is always Unpin, so is this type.
|
||||||
impl<T> Unpin for ReusableBoxFuture<T> {}
|
impl<T> Unpin for ReusableBoxFuture<'_, T> {}
|
||||||
|
|
||||||
impl<T> Drop for ReusableBoxFuture<T> {
|
impl<T> Drop for ReusableBoxFuture<'_, T> {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
unsafe {
|
unsafe {
|
||||||
drop(Box::from_raw(self.boxed.as_ptr()));
|
drop(Box::from_raw(self.boxed.as_ptr()));
|
||||||
@@ -144,7 +141,7 @@ impl<T> Drop for ReusableBoxFuture<T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T> fmt::Debug for ReusableBoxFuture<T> {
|
impl<T> fmt::Debug for ReusableBoxFuture<'_, T> {
|
||||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||||
f.debug_struct("ReusableBoxFuture").finish()
|
f.debug_struct("ReusableBoxFuture").finish()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user