use std::collections::VecDeque;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Semaphore;
use tokio::sync::oneshot::Sender;
use tracing::trace;
use crate::client::conn::Connector;
use crate::client::conn::Protocol;
use crate::client::conn::Transport;
use super::IdleConnections;
use super::PoolableConnection;
use super::Pooled;
use super::lock::ArcMutex;
use super::lock::ArcMutexGuard;
use super::lock::WeakMutex;
pub use super::checkout::Checkout;
#[derive(Debug, Clone)]
#[cfg_attr(feature = "serde", derive(::serde::Serialize, ::serde::Deserialize))]
#[non_exhaustive]
pub struct ConnectionManagerConfig {
pub idle_timeout: Option<Duration>,
pub max_idle_per_host: usize,
pub continue_after_preemption: bool,
pub max_connecting_per_host: Option<usize>,
}
impl Default for ConnectionManagerConfig {
fn default() -> Self {
Self {
idle_timeout: Some(Duration::from_secs(90)),
max_idle_per_host: 32,
continue_after_preemption: true,
max_connecting_per_host: None,
}
}
}
#[derive(Debug)]
pub struct ConnectionManager<C, R>
where
C: PoolableConnection<R>,
R: Send + 'static,
{
inner: ArcMutex<InnerConnectionManager<C, R>>,
}
impl<C, R> Clone for ConnectionManager<C, R>
where
C: PoolableConnection<R>,
R: Send + 'static,
{
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
}
}
}
impl<C, R> ConnectionManager<C, R>
where
C: PoolableConnection<R>,
R: Send + 'static,
{
pub fn new(config: impl Into<Arc<ConnectionManagerConfig>>) -> Self {
Self {
inner: InnerConnectionManager::new(config),
}
}
pub fn checkout<T, P>(&self, connector: Connector<T, P, R>) -> Checkout<T, P, R>
where
T: Transport<R> + Send,
P: Protocol<T::IO, R, Connection = C> + Send + 'static,
{
InnerConnectionManager::checkout(&mut self.inner.lock(), connector)
}
}
#[derive(Debug)]
pub(super) struct InnerConnectionManager<C, R>
where
C: PoolableConnection<R>,
R: Send + 'static,
{
connecting: bool,
waiting: VecDeque<Sender<Pooled<C, R>>>,
idle: IdleConnections<C, R>,
config: Arc<ConnectionManagerConfig>,
connecting_permits: Option<Arc<Semaphore>>,
}
impl<C, R> InnerConnectionManager<C, R>
where
C: PoolableConnection<R>,
R: Send + 'static,
{
fn new(config: impl Into<Arc<ConnectionManagerConfig>>) -> ArcMutex<Self> {
let config = config.into();
let connecting_permits = config
.max_connecting_per_host
.map(|limit| Arc::new(Semaphore::new(limit)));
ArcMutex::new(Self {
connecting: false,
waiting: VecDeque::new(),
idle: IdleConnections::default(),
config,
connecting_permits,
})
}
fn checkout<T, P>(
manager: &mut ArcMutexGuard<Self>,
mut connector: Connector<T, P, R>,
) -> Checkout<T, P, R>
where
T: Transport<R> + Send,
P: Protocol<T::IO, R, Connection = C> + Send + 'static,
{
let (tx, rx) = tokio::sync::oneshot::channel();
let multiplex = connector.multiplex();
if let Some(connection) = manager.pop() {
trace!("connection found in pool");
let request = connector.take_request_unpinned();
return Checkout::new_connected(manager.downgrade(), rx, connection, request);
}
trace!("checkout interested in pooled connections");
manager.waiting.push_back(tx);
if manager.connecting {
trace!("connection in progress elsewhere, will wait");
Checkout::new_connecting(
manager.downgrade(),
rx,
connector,
manager.connecting_permits.clone(),
)
} else {
let is_leader = multiplex;
if multiplex {
trace!("checkout of multiplexed connection, other connections should wait");
manager.connecting = true;
}
trace!("connecting to host");
Checkout::new_idle(
manager.downgrade(),
rx,
connector,
&manager.config,
is_leader,
manager.connecting_permits.clone(),
)
}
}
pub(in crate::client) fn cancel_connection(&mut self) {
if self.connecting {
trace!("pending connection cancelled");
}
self.connecting = false;
}
pub(in crate::client) fn connection_failed(&mut self) {
while let Some(waiter) = self.waiting.pop_front() {
if waiter.is_closed() {
continue;
}
trace!("connection attempt failed, releasing next waiter in line");
return;
}
trace!("connection attempt failed, no other waiters remain");
self.connecting = false;
}
pub(in crate::client) fn connected_in_handshake(&mut self, multiplex: bool) {
self.connecting = multiplex;
if !multiplex {
trace!(waiters=%self.waiting.len(), "dropping waiters");
self.waiting.clear();
}
}
pub(super) fn push(&mut self, mut connection: C, manager: &WeakMutex<Self>) {
self.connecting = false;
let _span = tracing::trace_span!("manager::push").entered();
trace!(waiters=%self.waiting.len(), "walking waiters");
while let Some(waiter) = self.waiting.pop_front() {
if waiter.is_closed() {
trace!("skipping closed waiter");
continue;
}
if let Some(conn) = connection.reuse() {
trace!("re-usable connection will be sent to waiter");
let pooled = Pooled {
connection: Some(conn),
manager: manager.clone(),
};
if waiter.send(pooled).is_err() {
trace!("waiter closed, skipping");
continue;
};
} else {
trace!("connection not re-usable, but will be sent to waiter");
let pooled = Pooled {
connection: Some(connection),
manager: manager.clone(),
};
let Err(pooled) = waiter.send(pooled) else {
trace!("connection sent");
return;
};
trace!("waiter closed, continuing");
connection = pooled.take().unwrap();
}
}
trace!("push idle connection");
self.idle
.push(connection, Some(self.config.max_idle_per_host));
}
pub(super) fn pop(&mut self) -> Option<C> {
self.idle.pop(self.config.idle_timeout)
}
}
#[cfg(all(test, feature = "mock"))]
mod tests {
use std::time::Duration;
use futures::FutureExt as _;
use crate::client::conn::transport::mock::MockConnectionError;
use super::*;
use crate::client::conn::connector::Error;
use crate::client::conn::protocol::mock::{MockProtocol, MockRequest, MockSender};
use crate::client::conn::stream::mock::MockStream;
use crate::client::conn::transport::mock::MockTransport;
#[tokio::test]
async fn checkout_simple() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let conn = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
let cid = conn.id();
drop(conn);
let conn = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
assert_eq!(conn.id(), cid, "connection should be re-used");
conn.close();
drop(conn);
let c2 = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(c2.is_open());
assert_ne!(c2.id(), cid, "connection should not be re-used");
}
#[tokio::test]
async fn checkout_multiplex() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let conn = manager
.checkout(MockTransport::reusable().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
let cid = conn.id();
drop(conn);
let conn = manager
.checkout(MockTransport::reusable().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
assert_eq!(conn.id(), cid, "connection should be re-used");
conn.close();
drop(conn);
let conn = manager
.checkout(MockTransport::reusable().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
assert_ne!(conn.id(), cid, "connection should not be re-used");
}
#[tokio::test]
async fn checkout_multiplex_contended() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (tx, rx) = tokio::sync::oneshot::channel();
let mut checkout_a =
std::pin::pin!(manager.checkout(MockTransport::channel(rx).connector(MockRequest),));
assert!(futures::poll!(&mut checkout_a).is_pending());
let mut checkout_b =
std::pin::pin!(manager.checkout(MockTransport::reusable().connector(MockRequest),));
assert!(futures::poll!(&mut checkout_b).is_pending());
assert!(tx.send(MockStream::reusable()).is_ok());
assert!(futures::poll!(&mut checkout_b).is_pending());
let conn_a = checkout_a.await.unwrap();
assert!(conn_a.is_open());
let conn_b = checkout_b.await.unwrap();
assert!(conn_b.is_open());
assert_eq!(conn_b.id(), conn_a.id(), "connection should be re-used");
}
#[tokio::test]
async fn checkout_idle_returned() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let conn = MockSender::single();
let first_id = conn.id();
let checkout = manager.checkout(MockTransport::single().connector(MockRequest));
manager
.inner
.lock()
.push(conn, &WeakMutex::downgrade(&manager.inner));
let conn = checkout.now_or_never().unwrap().unwrap();
assert!(conn.is_open());
assert_eq!(conn.id(), first_id, "connection should be re-used");
}
#[tokio::test]
async fn checkout_idle_connected() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let conn_first = MockSender::single();
let first_id = conn_first.id();
tracing::debug!("Checkout from pool");
let checkout = manager.checkout(MockTransport::single().connector(MockRequest));
tracing::debug!("Checking interest");
assert!(
!manager.inner.lock().waiting.is_empty(),
"No connections are waiting"
);
tracing::debug!("Resolving checkout");
let conn = checkout.now_or_never().unwrap().unwrap();
tracing::debug!("Inserting original connection");
manager
.inner
.lock()
.push(conn_first, &WeakMutex::downgrade(&manager.inner));
assert!(conn.is_open());
assert_ne!(conn.id(), first_id, "connection should not be re-used");
}
#[tokio::test]
async fn checkout_drop_pool_recover() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let start = manager.checkout(MockTransport::reusable().connector(MockRequest));
let checkout = manager.checkout(MockTransport::reusable().connector(MockRequest));
drop(start);
drop(manager);
assert!(checkout.now_or_never().unwrap().is_ok());
}
#[tokio::test]
async fn checkout_drop_pool() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let checkout = manager.checkout(MockTransport::reusable().connector(MockRequest));
drop(manager);
assert!(checkout.now_or_never().unwrap().is_ok());
}
#[tokio::test]
async fn checkout_connection_error() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let checkout = manager.checkout(MockTransport::error().connector(MockRequest));
let outcome = checkout.now_or_never().unwrap();
let error = outcome.unwrap_err();
assert!(matches!(error, Error::Connecting(MockConnectionError)));
}
#[tokio::test]
async fn checkout_pool_cloned() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let other = manager.clone();
let conn = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
let cid = conn.id();
drop(conn);
let conn = other
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
assert_eq!(conn.id(), cid, "connection should be re-used");
conn.close();
drop(conn);
let c2 = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(c2.is_open());
assert_ne!(c2.id(), cid, "connection should not be re-used");
}
#[tokio::test]
async fn checkout_delayed_drop() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: true,
..Default::default()
});
let conn = manager
.checkout(MockTransport::single().connector(MockRequest))
.await
.unwrap();
assert!(conn.is_open());
let cid = conn.id();
let checkout = manager.checkout(MockTransport::single().connector(MockRequest));
drop(conn);
let conn = checkout.await.unwrap();
assert!(conn.is_open());
assert_eq!(cid, conn.id());
assert_eq!(manager.inner.lock().idle.len(), 1);
}
#[tokio::test]
async fn checkout_connection_failure_releases_waiters() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (stream_tx, stream_rx) = tokio::sync::oneshot::channel();
let leader = manager.checkout(MockTransport::channel(stream_rx).connector(MockRequest));
let follower = manager.checkout(MockTransport::reusable().connector(MockRequest));
assert_eq!(
manager.inner.lock().waiting.len(),
2,
"leader and follower should both be queued waiting"
);
drop(stream_tx);
let leader_result = leader.await;
assert!(leader_result.is_err(), "leader should fail to connect");
let follower_result = tokio::time::timeout(Duration::from_secs(1), follower)
.await
.expect("follower should not hang waiting on the failed leader");
assert!(follower_result.unwrap().is_open());
}
#[tokio::test]
async fn checkout_connection_failure_wakes_waiting_checkout() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (stream_tx, stream_rx) = tokio::sync::oneshot::channel();
let leader = manager.checkout(MockTransport::channel(stream_rx).connector(MockRequest));
let follower = manager.checkout(MockTransport::reusable().connector(MockRequest));
let follower_task = tokio::spawn(follower);
tokio::task::yield_now().await;
tokio::task::yield_now().await;
drop(stream_tx);
let leader_result = leader.await;
assert!(leader_result.is_err(), "leader should fail to connect");
let follower_result = tokio::time::timeout(Duration::from_secs(1), follower_task)
.await
.expect("follower task should be woken and complete promptly")
.expect("follower task should not panic");
assert!(follower_result.unwrap().is_open());
}
#[tokio::test]
async fn checkout_connection_failure_releases_one_waiter_at_a_time() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (leader_tx, leader_rx) = tokio::sync::oneshot::channel();
let (follower_a_tx, follower_a_rx) = tokio::sync::oneshot::channel();
let mut leader = std::pin::pin!(
manager.checkout(MockTransport::channel(leader_rx).connector(MockRequest))
);
assert!(futures::poll!(&mut leader).is_pending());
let mut follower_a = std::pin::pin!(
manager.checkout(MockTransport::channel(follower_a_rx).connector(MockRequest))
);
assert!(futures::poll!(&mut follower_a).is_pending());
let mut follower_b =
std::pin::pin!(manager.checkout(MockTransport::reusable().connector(MockRequest)));
assert!(futures::poll!(&mut follower_b).is_pending());
assert_eq!(
manager.inner.lock().waiting.len(),
3,
"leader and both followers should be queued waiting"
);
drop(leader_tx);
let leader_result = leader.await;
assert!(leader_result.is_err(), "leader should fail to connect");
assert_eq!(
manager.inner.lock().waiting.len(),
1,
"only follower_b should remain queued after the leader fails"
);
assert!(futures::poll!(&mut follower_b).is_pending());
assert!(futures::poll!(&mut follower_a).is_pending());
drop(follower_a_tx);
let follower_a_result = follower_a.await;
assert!(
follower_a_result.is_err(),
"follower_a should fail to connect"
);
assert_eq!(
manager.inner.lock().waiting.len(),
0,
"follower_b should have been released after follower_a also failed"
);
let follower_b_result = tokio::time::timeout(Duration::from_secs(1), follower_b)
.await
.expect("follower_b should not hang once released");
assert!(follower_b_result.unwrap().is_open());
}
#[tokio::test]
async fn checkout_multiplex_ready_false_releases_waiters() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let leader_connector = Connector::new(
MockTransport::reusable(),
MockProtocol::new(true).with_multiplex_ready(false),
MockRequest,
);
let leader = manager.checkout(leader_connector);
let follower = manager.checkout(MockTransport::reusable().connector(MockRequest));
assert_eq!(
manager.inner.lock().waiting.len(),
2,
"leader and follower should both be queued waiting"
);
let leader_conn = leader.await.unwrap();
assert!(leader_conn.is_open());
let follower_conn = tokio::time::timeout(Duration::from_secs(1), follower)
.await
.expect("follower should not hang waiting on a non-multiplexed leader")
.unwrap();
assert!(follower_conn.is_open());
assert_ne!(
follower_conn.id(),
leader_conn.id(),
"follower should have connected on its own, not shared the leader's connection"
);
}
#[tokio::test]
async fn checkout_leader_abandoned_releases_waiters() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (_stream_tx, stream_rx) = tokio::sync::oneshot::channel();
let leader = manager.checkout(MockTransport::channel(stream_rx).connector(MockRequest));
let follower = manager.checkout(MockTransport::reusable().connector(MockRequest));
assert_eq!(
manager.inner.lock().waiting.len(),
2,
"leader and follower should both be queued waiting"
);
drop(leader);
let follower_result = tokio::time::timeout(Duration::from_secs(1), follower)
.await
.expect("follower should not hang waiting on an abandoned leader");
assert!(follower_result.unwrap().is_open());
assert!(
!manager.inner.lock().connecting,
"connecting flag should be reset"
);
}
#[tokio::test]
async fn checkout_follower_drop_does_not_disrupt_leader_or_other_waiters() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
..Default::default()
});
let (stream_tx, stream_rx) = tokio::sync::oneshot::channel();
let mut leader = std::pin::pin!(
manager.checkout(MockTransport::channel(stream_rx).connector(MockRequest))
);
assert!(futures::poll!(&mut leader).is_pending());
let follower_a = manager.checkout(MockTransport::reusable().connector(MockRequest));
let follower_b = manager.checkout(MockTransport::reusable().connector(MockRequest));
drop(follower_a);
assert!(stream_tx.send(MockStream::reusable()).is_ok());
let leader_conn = leader.await.unwrap();
assert!(leader_conn.is_open());
let follower_b_conn = tokio::time::timeout(Duration::from_secs(1), follower_b)
.await
.expect("follower_b should not hang")
.unwrap();
assert_eq!(
follower_b_conn.id(),
leader_conn.id(),
"follower_b should still share the leader's connection despite follower_a's drop"
);
}
#[tokio::test]
async fn checkout_limits_simultaneous_connection_attempts() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
max_connecting_per_host: Some(1),
..Default::default()
});
let (tx_a, rx_a) = tokio::sync::oneshot::channel();
let connector_a = Connector::new(
MockTransport::channel(rx_a),
MockProtocol::new(false),
MockRequest,
);
let connector_b = Connector::new(
MockTransport::single(),
MockProtocol::new(false),
MockRequest,
);
let mut checkout_a = std::pin::pin!(manager.checkout(connector_a));
let mut checkout_b = std::pin::pin!(manager.checkout(connector_b));
assert!(futures::poll!(&mut checkout_a).is_pending());
assert!(futures::poll!(&mut checkout_b).is_pending());
assert!(tx_a.send(MockStream::single()).is_ok());
let conn_a = checkout_a.await.unwrap();
assert!(conn_a.is_open());
let conn_b = tokio::time::timeout(Duration::from_secs(1), checkout_b)
.await
.expect("checkout_b should proceed once a permit frees up")
.unwrap();
assert!(conn_b.is_open());
}
#[tokio::test]
async fn checkout_connecting_limit_throttles_full_queue_release() {
let _ = tracing_subscriber::fmt::try_init();
let manager = ConnectionManager::<MockSender, MockRequest>::new(ConnectionManagerConfig {
idle_timeout: Some(Duration::from_secs(10)),
max_idle_per_host: 5,
continue_after_preemption: false,
max_connecting_per_host: Some(1),
..Default::default()
});
let leader_connector = Connector::new(
MockTransport::reusable(),
MockProtocol::new(true).with_multiplex_ready(false),
MockRequest,
);
let leader = manager.checkout(leader_connector);
let (tx_a, rx_a) = tokio::sync::oneshot::channel();
let (tx_b, rx_b) = tokio::sync::oneshot::channel();
let follower_a_connector = Connector::new(
MockTransport::channel(rx_a),
MockProtocol::new(false),
MockRequest,
);
let follower_b_connector = Connector::new(
MockTransport::channel(rx_b),
MockProtocol::new(false),
MockRequest,
);
let mut follower_a = std::pin::pin!(manager.checkout(follower_a_connector));
let mut follower_b = std::pin::pin!(manager.checkout(follower_b_connector));
assert!(futures::poll!(&mut follower_a).is_pending());
assert!(futures::poll!(&mut follower_b).is_pending());
let leader_conn = leader.await.unwrap();
assert!(leader_conn.is_open());
assert!(futures::poll!(&mut follower_a).is_pending());
assert!(futures::poll!(&mut follower_b).is_pending());
assert!(tx_a.send(MockStream::single()).is_ok());
let conn_a = follower_a.await.unwrap();
assert!(conn_a.is_open());
assert!(tx_b.send(MockStream::single()).is_ok());
let conn_b = tokio::time::timeout(Duration::from_secs(1), follower_b)
.await
.expect("follower_b should proceed once a permit frees up")
.unwrap();
assert!(conn_b.is_open());
}
}