use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc;
use std::task::{Context, Poll, Waker};
use std::time::{Duration, Instant};
use parking_lot::Mutex as PoolMutex;
use smallvec::SmallVec;
use crate::cx::Cx;
use super::waiter::DeferredWaker;
fn wall_clock_now() -> Instant {
Instant::now()
}
pub type PoolFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
pub type PoolReturnSender<R> = mpsc::Sender<PoolReturn<R>>;
pub type PoolReturnReceiver<R> = mpsc::Receiver<PoolReturn<R>>;
type ReturnWakerEntry = (u64, DeferredWaker);
type ReturnWakerList = SmallVec<[ReturnWakerEntry; 4]>;
type ReturnWakers = Arc<PoolMutex<ReturnWakerList>>;
pub trait Pool: Send + Sync {
type Resource: Send;
type Error: std::error::Error + Send + Sync + 'static;
fn acquire<'a>(
&'a self,
cx: &'a Cx,
) -> PoolFuture<'a, Result<PooledResource<Self::Resource>, Self::Error>>;
fn try_acquire(&self) -> Option<PooledResource<Self::Resource>>;
fn stats(&self) -> PoolStats;
fn close(&self) -> PoolFuture<'_, ()>;
fn health_check<'a>(&'a self, _resource: &'a Self::Resource) -> PoolFuture<'a, bool> {
Box::pin(async { true })
}
}
pub trait AsyncResourceFactory: Send + Sync {
type Resource: Send;
type Error: Send + Sync + 'static + Into<Box<dyn std::error::Error + Send + Sync>>;
#[allow(clippy::type_complexity)]
fn create(
&self,
) -> Pin<Box<dyn Future<Output = Result<Self::Resource, Self::Error>> + Send + '_>>;
}
#[derive(Debug, Clone, Default)]
pub struct PoolStats {
pub active: usize,
pub idle: usize,
pub total: usize,
pub max_size: usize,
pub waiters: usize,
pub total_acquisitions: u64,
pub total_wait_time: Duration,
}
#[derive(Debug)]
pub enum PoolReturn<R> {
Return {
resource: R,
hold_duration: Duration,
created_at: Instant,
},
Discard {
hold_duration: Duration,
},
}
#[derive(Debug)]
struct ReturnObligation {
discharged: bool,
}
impl ReturnObligation {
#[inline]
fn new() -> Self {
Self { discharged: false }
}
#[inline]
fn discharge(&mut self) {
self.discharged = true;
}
#[inline]
fn is_discharged(&self) -> bool {
self.discharged
}
}
#[must_use = "PooledResource must be returned or dropped"]
pub struct PooledResource<R> {
resource: Option<R>,
return_obligation: ReturnObligation,
return_tx: PoolReturnSender<R>,
acquired_at: Instant,
created_at: Instant,
time_getter: fn() -> Instant,
return_wakers: Option<ReturnWakers>,
#[allow(dead_code)]
is_broken: bool,
}
impl<R> PooledResource<R> {
#[inline]
pub fn new(resource: R, return_tx: PoolReturnSender<R>) -> Self {
Self::new_with_time_getter(resource, return_tx, wall_clock_now)
}
#[inline]
pub fn new_with_time_getter(
resource: R,
return_tx: PoolReturnSender<R>,
time_getter: fn() -> Instant,
) -> Self {
let now = time_getter();
Self::new_with_timestamps(resource, return_tx, now, now, time_getter)
}
fn new_with_timestamps(
resource: R,
return_tx: PoolReturnSender<R>,
acquired_at: Instant,
created_at: Instant,
time_getter: fn() -> Instant,
) -> Self {
Self {
resource: Some(resource),
return_obligation: ReturnObligation::new(),
return_tx,
acquired_at,
created_at,
time_getter,
return_wakers: None,
is_broken: false,
}
}
fn with_return_notify(mut self, wakers: ReturnWakers) -> Self {
self.return_wakers = Some(wakers);
self
}
#[inline]
#[must_use]
pub fn get(&self) -> &R {
self.resource.as_ref().expect(
"PooledResource accessed after drop or return - resource has been taken. \
This indicates a use-after-drop bug or concurrent access violation.",
)
}
#[inline]
pub fn get_mut(&mut self) -> &mut R {
self.resource.as_mut().expect(
"PooledResource accessed after drop or return - resource has been taken. \
This indicates a use-after-drop bug or concurrent access violation.",
)
}
#[inline]
#[must_use]
pub fn try_get(&self) -> Option<&R> {
self.resource.as_ref()
}
#[inline]
pub fn try_get_mut(&mut self) -> Option<&mut R> {
self.resource.as_mut()
}
pub fn return_to_pool(mut self) {
self.return_inner();
}
pub fn discard(mut self) {
self.discard_inner();
}
#[inline]
pub fn mark_broken(&mut self) {
self.is_broken = true;
}
#[inline]
#[must_use]
pub fn is_broken(&self) -> bool {
self.is_broken
}
#[inline]
#[must_use]
pub fn held_duration(&self) -> Duration {
(self.time_getter)().saturating_duration_since(self.acquired_at)
}
fn return_inner(&mut self) {
if self.return_obligation.is_discharged() {
return;
}
let hold_duration = self.held_duration();
if let Some(resource) = self.resource.take() {
let _ = self.return_tx.send(PoolReturn::Return {
resource,
hold_duration,
created_at: self.created_at,
});
}
self.return_obligation.discharge();
self.notify_return_wakers();
}
fn discard_inner(&mut self) {
if self.return_obligation.is_discharged() {
return;
}
let hold_duration = self.held_duration();
let discarded = self.resource.take();
let _ = self.return_tx.send(PoolReturn::Discard { hold_duration });
self.return_obligation.discharge();
self.notify_return_wakers();
drop(discarded);
}
fn notify_return_wakers(&self) {
if let Some(ref wakers) = self.return_wakers {
let waker = {
let lock = wakers.lock();
lock.first().map(|(_, waker)| waker.clone_waker())
};
if let Some(waker) = waker {
waker.wake();
}
}
}
}
impl<R> Drop for PooledResource<R> {
fn drop(&mut self) {
if self.is_broken {
self.discard_inner();
} else {
self.return_inner();
}
}
}
impl<R> std::ops::Deref for PooledResource<R> {
type Target = R;
#[inline]
fn deref(&self) -> &Self::Target {
self.get()
}
}
impl<R> std::ops::DerefMut for PooledResource<R> {
#[inline]
fn deref_mut(&mut self) -> &mut Self::Target {
self.get_mut()
}
}
#[must_use = "a checked pool resource must be returned or dropped"]
pub struct CheckedPooledResource<R> {
pooled: Option<PooledResource<R>>,
obligation: Option<crate::runtime::obligation_mailbox::ObligationToken>,
}
impl<R> CheckedPooledResource<R> {
fn admit(pooled: PooledResource<R>, cx: &Cx) -> Result<Self, CheckedPoolError> {
let mut checked = Self {
pooled: Some(pooled),
obligation: None,
};
checked.obligation =
cx.try_register_obligation_checked(crate::record::ObligationKind::Lease, cx.task_id())?;
Ok(checked)
}
#[must_use]
pub fn get(&self) -> &R {
self.pooled.as_ref().expect("active checked checkout").get()
}
pub fn get_mut(&mut self) -> &mut R {
self.pooled
.as_mut()
.expect("active checked checkout")
.get_mut()
}
pub fn mark_broken(&mut self) {
self.pooled
.as_mut()
.expect("active checked checkout")
.mark_broken();
}
#[must_use]
pub fn is_broken(&self) -> bool {
self.pooled
.as_ref()
.expect("active checked checkout")
.is_broken()
}
#[must_use]
pub fn held_duration(&self) -> Duration {
self.pooled
.as_ref()
.expect("active checked checkout")
.held_duration()
}
pub fn return_to_pool(mut self) {
self.finish(true, false);
}
pub fn discard(mut self) {
self.finish(false, true);
}
fn finish(&mut self, explicit_return: bool, discard: bool) {
let Some(mut pooled) = self.pooled.take() else {
return;
};
pooled.return_obligation.discharge();
let mut resource = pooled.resource.take();
let discard = discard || pooled.is_broken;
let mut panics = CheckedPoolCleanupPanics::new();
let mut hold_duration = Duration::ZERO;
panics.run(|| hold_duration = pooled.held_duration());
let message = if discard {
PoolReturn::Discard { hold_duration }
} else {
PoolReturn::Return {
resource: resource.take().expect("active checked resource"),
hold_duration,
created_at: pooled.created_at,
}
};
let returned = match pooled.return_tx.send(message) {
Ok(()) => true,
Err(mpsc::SendError(PoolReturn::Return { resource: lost, .. })) => {
resource = Some(lost);
false
}
Err(mpsc::SendError(PoolReturn::Discard { .. })) => false,
};
let notification = self.obligation.take().and_then(|token| {
let (_, notification) = if explicit_return && !discard && returned {
token.commit_deferred()
} else {
token.abort_deferred(crate::record::ObligationAbortReason::Explicit)
};
notification
});
let waker = pooled
.return_wakers
.as_ref()
.and_then(|wakers| wakers.lock().first().map(|(_, waker)| waker.clone_waker()));
if let Some(gateway) = notification {
panics.run(|| gateway.notify());
panics.run(|| drop(gateway));
}
if let Some(waker) = waker {
panics.run(|| waker.wake_by_ref());
panics.run(|| drop(waker));
}
panics.run(|| drop(resource));
panics.run(|| drop(pooled));
}
}
impl<R> Drop for CheckedPooledResource<R> {
fn drop(&mut self) {
self.finish(false, false);
}
}
impl<R> std::ops::Deref for CheckedPooledResource<R> {
type Target = R;
fn deref(&self) -> &R {
self.get()
}
}
impl<R> std::ops::DerefMut for CheckedPooledResource<R> {
fn deref_mut(&mut self) -> &mut R {
self.get_mut()
}
}
impl<R: std::fmt::Debug> std::fmt::Debug for CheckedPooledResource<R> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CheckedPooledResource")
.field(
"resource",
&self.pooled.as_ref().and_then(PooledResource::try_get),
)
.field("tracked", &self.obligation.is_some())
.finish_non_exhaustive()
}
}
struct CheckedPoolCleanupPanics {
already_unwinding: bool,
first: Option<Box<dyn std::any::Any + Send>>,
}
impl CheckedPoolCleanupPanics {
fn new() -> Self {
Self {
already_unwinding: std::thread::panicking(),
first: None,
}
}
fn run(&mut self, callback: impl FnOnce()) {
if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(callback)) {
if self.already_unwinding || self.first.is_some() {
std::mem::forget(payload);
} else {
self.first = Some(payload);
}
}
}
}
impl Drop for CheckedPoolCleanupPanics {
fn drop(&mut self) {
if let Some(payload) = self.first.take() {
if std::thread::panicking() {
std::mem::forget(payload);
} else {
std::panic::resume_unwind(payload);
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum WarmupStrategy {
#[default]
BestEffort,
FailFast,
RequireMinimum,
}
#[derive(Debug, Clone)]
pub struct PoolConfig {
pub min_size: usize,
pub max_size: usize,
pub acquire_timeout: Duration,
pub idle_timeout: Duration,
pub max_lifetime: Duration,
pub health_check_on_acquire: bool,
pub health_check_interval: Option<Duration>,
pub evict_unhealthy: bool,
pub warmup_connections: usize,
pub warmup_timeout: Duration,
pub warmup_failure_strategy: WarmupStrategy,
}
impl Default for PoolConfig {
fn default() -> Self {
Self {
min_size: 1,
max_size: 10,
acquire_timeout: Duration::from_secs(30),
idle_timeout: Duration::from_mins(10),
max_lifetime: Duration::from_hours(1),
health_check_on_acquire: false,
health_check_interval: None,
evict_unhealthy: true,
warmup_connections: 0,
warmup_timeout: Duration::from_secs(30),
warmup_failure_strategy: WarmupStrategy::BestEffort,
}
}
}
impl PoolConfig {
#[must_use]
pub fn with_max_size(max_size: usize) -> Self {
Self {
max_size,
..Default::default()
}
}
#[must_use]
pub fn min_size(mut self, min_size: usize) -> Self {
self.min_size = min_size;
self
}
#[must_use]
pub fn max_size(mut self, max_size: usize) -> Self {
self.max_size = max_size;
self
}
#[must_use]
pub fn acquire_timeout(mut self, timeout: Duration) -> Self {
self.acquire_timeout = timeout;
self
}
#[must_use]
pub fn idle_timeout(mut self, timeout: Duration) -> Self {
self.idle_timeout = timeout;
self
}
#[must_use]
pub fn max_lifetime(mut self, lifetime: Duration) -> Self {
self.max_lifetime = lifetime;
self
}
#[must_use]
pub fn health_check_on_acquire(mut self, enabled: bool) -> Self {
self.health_check_on_acquire = enabled;
self
}
#[must_use]
pub fn health_check_interval(mut self, interval: Option<Duration>) -> Self {
self.health_check_interval = interval;
self
}
#[must_use]
pub fn evict_unhealthy(mut self, evict: bool) -> Self {
self.evict_unhealthy = evict;
self
}
#[must_use]
pub fn warmup_connections(mut self, count: usize) -> Self {
self.warmup_connections = count;
self
}
#[must_use]
pub fn warmup_timeout(mut self, timeout: Duration) -> Self {
self.warmup_timeout = timeout;
self
}
#[must_use]
pub fn warmup_failure_strategy(mut self, strategy: WarmupStrategy) -> Self {
self.warmup_failure_strategy = strategy;
self
}
}
#[derive(Debug)]
pub enum PoolError {
Closed,
Timeout,
Cancelled,
CreateFailed(Box<dyn std::error::Error + Send + Sync>),
}
impl std::fmt::Display for PoolError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Closed => write!(f, "pool closed"),
Self::Timeout => write!(f, "pool acquire timeout"),
Self::Cancelled => write!(f, "pool acquire cancelled"),
Self::CreateFailed(e) => write!(f, "resource creation failed: {e}"),
}
}
}
impl std::error::Error for PoolError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::CreateFailed(e) => Some(e.as_ref()),
_ => None,
}
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum CheckedPoolError {
#[error(transparent)]
Pool(#[from] PoolError),
#[error(transparent)]
Admission(#[from] crate::runtime::obligation_mailbox::ObligationAdmissionError),
}
#[derive(Debug)]
struct IdleResource<R> {
resource: R,
idle_since: Instant,
created_at: Instant,
}
struct PoolWaiter {
id: u64,
waker: DeferredWaker,
}
struct DetachedPoolWakers {
state_waker: Option<DeferredWaker>,
return_waker: Option<DeferredWaker>,
next_dispatcher: Option<Waker>,
}
impl DetachedPoolWakers {
fn retire(self) {
let Self {
state_waker,
return_waker,
next_dispatcher,
} = self;
if let Some(next) = next_dispatcher {
next.wake();
}
drop(state_waker);
drop(return_waker);
}
}
struct GenericPoolState<R> {
idle: std::collections::VecDeque<IdleResource<R>>,
active: usize,
creating: usize,
total_created: u64,
total_acquisitions: u64,
total_wait_time: Duration,
closed: bool,
waiters: std::collections::VecDeque<PoolWaiter>,
next_waiter_id: u64,
}
struct WaitForNotification<'a, 'b, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
pool: &'a GenericPool<R, F>,
waiter_id: &'b mut Option<u64>,
cx: &'a Cx,
completed: bool,
}
impl<R, F> Future for WaitForNotification<'_, '_, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
type Output = ();
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
if self.cx.checkpoint().is_err() {
self.completed = true;
return Poll::Ready(());
}
self.pool.process_returns();
let mut state_candidate = Some(DeferredWaker::new(cx.waker().clone()));
let mut return_candidate = Some(DeferredWaker::new(cx.waker().clone()));
let mut retired_state_waker = None;
let mut state = self.pool.state.lock();
if state.closed {
self.completed = true;
return Poll::Ready(());
}
let total_including_creating = state.active + state.idle.len() + state.creating;
let available = state.idle.len()
+ self
.pool
.config
.max_size
.saturating_sub(total_including_creating);
let pos = if let Some(id) = *self.waiter_id {
if let Some(idx) = state.waiters.iter().position(|w| w.id == id) {
if idx >= available {
let w = &mut state.waiters[idx];
let candidate = state_candidate
.take()
.expect("state waker candidate is available");
if w.waker.will_wake(cx.waker()) {
retired_state_waker = Some(candidate);
} else {
retired_state_waker = Some(std::mem::replace(&mut w.waker, candidate));
}
}
idx
} else {
state.waiters.reserve(1);
state.waiters.push_front(PoolWaiter {
id,
waker: state_candidate
.take()
.expect("state waker candidate is available"),
});
0
}
} else {
let id = state.next_waiter_id;
state.next_waiter_id = state.next_waiter_id.wrapping_add(1);
let idx = state.waiters.len();
state.waiters.reserve(1);
state.waiters.push_back(PoolWaiter {
id,
waker: state_candidate
.take()
.expect("state waker candidate is available"),
});
*self.waiter_id = Some(id);
idx
};
let id = self.waiter_id.expect("waiter_id assigned above");
drop(state);
if pos < available {
self.completed = true;
drop(retired_state_waker);
drop(state_candidate);
drop(return_candidate);
return Poll::Ready(());
}
let mut retired_return_waker = None;
{
let mut wakers = self.pool.return_wakers.lock();
if let Some((_, existing)) = wakers.iter_mut().find(|(wid, _)| *wid == id) {
let candidate = return_candidate
.take()
.expect("return waker candidate is available");
if existing.will_wake(cx.waker()) {
retired_return_waker = Some(candidate);
} else {
retired_return_waker = Some(std::mem::replace(existing, candidate));
}
} else {
wakers.reserve(1);
wakers.push((
id,
return_candidate
.take()
.expect("return waker candidate is available"),
));
}
}
let closed_after_registration = self.pool.closed.load(Ordering::Acquire);
drop(retired_state_waker);
drop(retired_return_waker);
drop(state_candidate);
drop(return_candidate);
if closed_after_registration {
self.completed = true;
return Poll::Ready(());
}
self.pool.process_returns();
Poll::Pending
}
}
impl<R, F> Drop for WaitForNotification<'_, '_, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn drop(&mut self) {
if !self.completed
&& let Some(id) = *self.waiter_id
{
let mut retired_state_waker = None;
let marginal_waker: Option<Waker> = {
let mut state = self.pool.state.lock();
if let Some(p) = state.waiters.iter().position(|w| w.id == id) {
retired_state_waker = state.waiters.remove(p).map(|waiter| waiter.waker);
if state.closed {
None
} else {
let total_including_creating =
state.active + state.idle.len() + state.creating;
let available = state.idle.len()
+ self
.pool
.config
.max_size
.saturating_sub(total_including_creating);
if p < available && available > 0 && available - 1 < state.waiters.len() {
Some(state.waiters[available - 1].waker.clone_waker())
} else {
None
}
}
} else {
None
}
};
let (retired_return_waker, next_dispatcher) = self.pool.detach_return_waker(id);
if let Some(waker) = marginal_waker {
waker.wake();
}
DetachedPoolWakers {
state_waker: retired_state_waker,
return_waker: retired_return_waker,
next_dispatcher,
}
.retire();
}
}
}
struct ActiveCheckoutGuard<'a, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
pool: &'a GenericPool<R, F>,
completed: bool,
health_check_in_progress: bool,
}
impl<'a, R, F> ActiveCheckoutGuard<'a, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn new(pool: &'a GenericPool<R, F>) -> Self {
Self {
pool,
completed: false,
health_check_in_progress: false,
}
}
fn begin_health_check(&mut self) {
self.health_check_in_progress = true;
}
fn finish_health_check(&mut self) {
self.health_check_in_progress = false;
}
fn reject_unhealthy(mut self) {
self.completed = true;
self.pool.reject_unhealthy_idle_resource();
}
fn commit(mut self) {
self.completed = true;
}
}
impl<R, F> Drop for ActiveCheckoutGuard<'_, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn drop(&mut self) {
if !self.completed {
if self.health_check_in_progress {
self.pool.reject_unhealthy_idle_resource();
} else {
self.pool.rollback_active_checkout();
}
}
}
}
struct CreateSlotClaim {
retired_waker: Option<DeferredWaker>,
}
struct CreateSlotReservation<'a, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
pool: &'a GenericPool<R, F>,
committed: bool,
}
impl<'a, R, F> CreateSlotReservation<'a, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn try_reserve(pool: &'a GenericPool<R, F>, waiter_id: Option<u64>) -> Option<Self> {
let CreateSlotClaim { retired_waker } = pool.reserve_create_slot(waiter_id)?;
let reservation = Self {
pool,
committed: false,
};
drop(retired_waker);
Some(reservation)
}
fn commit(mut self) -> bool {
let handed_out = self.pool.commit_create_slot();
self.committed = true;
handed_out
}
fn committed_manually(mut self) {
self.committed = true;
}
}
impl<R, F> Drop for CreateSlotReservation<'_, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn drop(&mut self) {
if !self.committed {
self.pool.release_create_slot();
}
}
}
impl<F, R, E, Fut> AsyncResourceFactory for F
where
F: Fn() -> Fut + Send + Sync,
Fut: Future<Output = Result<R, E>> + Send + 'static,
R: Send,
E: Send + Sync + 'static + Into<Box<dyn std::error::Error + Send + Sync>>,
{
type Resource = R;
type Error = E;
fn create(
&self,
) -> Pin<Box<dyn Future<Output = Result<Self::Resource, Self::Error>> + Send + '_>> {
Box::pin(self())
}
}
pub struct GenericPool<R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
factory: F,
config: PoolConfig,
state: PoolMutex<GenericPoolState<R>>,
return_tx: PoolReturnSender<R>,
return_rx: PoolMutex<PoolReturnReceiver<R>>,
time_getter: fn() -> Instant,
#[allow(clippy::type_complexity)]
health_check_fn: Option<Box<dyn Fn(&R) -> bool + Send + Sync>>,
return_wakers: ReturnWakers,
closed: AtomicBool,
#[cfg(feature = "metrics")]
metrics: Option<PoolMetricsHandle>,
}
impl<R, F> GenericPool<R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
pub fn new(factory: F, config: PoolConfig) -> Self {
Self::with_time_getter(factory, config, wall_clock_now)
}
pub fn with_time_getter(factory: F, config: PoolConfig, time_getter: fn() -> Instant) -> Self {
let (return_tx, return_rx) = mpsc::channel();
Self {
factory,
config,
state: PoolMutex::new(GenericPoolState {
idle: std::collections::VecDeque::with_capacity(8),
active: 0,
creating: 0,
total_created: 0,
total_acquisitions: 0,
total_wait_time: Duration::ZERO,
closed: false,
waiters: std::collections::VecDeque::with_capacity(4),
next_waiter_id: 0,
}),
return_tx,
return_rx: PoolMutex::new(return_rx),
time_getter,
health_check_fn: None,
return_wakers: Arc::new(PoolMutex::new(SmallVec::new())),
closed: AtomicBool::new(false),
#[cfg(feature = "metrics")]
metrics: None,
}
}
pub fn with_factory(factory: F) -> Self {
Self::new(factory, PoolConfig::default())
}
#[must_use]
pub const fn time_getter(&self) -> fn() -> Instant {
self.time_getter
}
pub fn acquire_checked<'a>(
&'a self,
cx: &'a Cx,
) -> PoolFuture<'a, Result<CheckedPooledResource<R>, CheckedPoolError>> {
Box::pin(async move {
let pooled = self.acquire(cx).await?;
CheckedPooledResource::admit(pooled, cx)
})
}
pub fn try_acquire_checked(
&self,
cx: &Cx,
) -> Result<Option<CheckedPooledResource<R>>, CheckedPoolError> {
cx.checkpoint().map_err(|_| PoolError::Cancelled)?;
if self.closed.load(Ordering::Acquire) {
return Err(PoolError::Closed.into());
}
match self.try_acquire() {
Some(pooled) => CheckedPooledResource::admit(pooled, cx).map(Some),
None if self.closed.load(Ordering::Acquire) => Err(PoolError::Closed.into()),
None => Ok(None),
}
}
#[cfg(feature = "metrics")]
#[must_use]
pub fn with_metrics(mut self, handle: PoolMetricsHandle) -> Self {
self.metrics = Some(handle);
self
}
#[must_use]
pub fn with_health_check(mut self, check: impl Fn(&R) -> bool + Send + Sync + 'static) -> Self {
self.health_check_fn = Some(Box::new(check));
self
}
pub async fn warmup(&self) -> Result<usize, PoolError> {
let mut created = 0;
let mut last_error = None;
let warmup_deadline = crate::time::wall_now() + self.config.warmup_timeout;
let target = self.config.warmup_connections;
for _ in 0..target {
let Some(slot) = CreateSlotReservation::try_reserve(self, None) else {
break; };
let now = crate::time::wall_now();
let remaining = Duration::from_nanos(warmup_deadline.duration_since(now));
let create_result = if remaining.is_zero() {
Err(PoolError::Timeout)
} else {
match crate::time::timeout(now, remaining, self.create_resource()).await {
Ok(result) => result,
Err(_elapsed) => Err(PoolError::Timeout),
}
};
match create_result {
Ok(resource) => {
slot.committed_manually();
self.commit_create_slot_as_idle(resource);
created += 1;
}
Err(PoolError::Timeout) => match self.config.warmup_failure_strategy {
WarmupStrategy::FailFast => return Err(PoolError::Timeout),
WarmupStrategy::BestEffort | WarmupStrategy::RequireMinimum => {
last_error = Some(PoolError::Timeout);
break;
}
},
Err(e) => {
match self.config.warmup_failure_strategy {
WarmupStrategy::FailFast => return Err(e),
WarmupStrategy::BestEffort | WarmupStrategy::RequireMinimum => {
last_error = Some(e);
}
}
}
}
}
if self.config.warmup_failure_strategy == WarmupStrategy::RequireMinimum
&& created < self.config.min_size
{
return Err(last_error.unwrap_or(PoolError::CreateFailed(
"warmup did not reach min_size".into(),
)));
}
Ok(created)
}
fn is_healthy(&self, resource: &R) -> bool {
self.health_check_fn
.as_ref()
.is_none_or(|check| check(resource))
}
fn rollback_active_checkout(&self) {
let waker = {
let mut state = self.state.lock();
state.active = state.active.saturating_sub(1);
state.total_acquisitions = state.total_acquisitions.saturating_sub(1);
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
if available > 0 && available - 1 < state.waiters.len() {
Some(state.waiters[available - 1].waker.clone_waker())
} else {
None
}
};
#[cfg(feature = "metrics")]
self.update_metrics_gauges();
if let Some(waker) = waker {
waker.wake();
}
}
fn reject_unhealthy_idle_resource(&self) {
self.rollback_active_checkout();
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_destroyed(DestroyReason::Unhealthy);
}
}
fn try_recv_return(&self) -> Result<PoolReturn<R>, mpsc::TryRecvError> {
let rx = self.return_rx.lock();
rx.try_recv()
}
#[cfg_attr(not(feature = "metrics"), allow(unused_variables))]
fn process_returns(&self) {
let mut waiters_to_wake: SmallVec<[Waker; 4]> = SmallVec::new();
while let Ok(ret) = self.try_recv_return() {
match ret {
PoolReturn::Return {
resource,
hold_duration,
created_at,
} => {
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_released(hold_duration);
}
#[allow(clippy::needless_late_init)]
let mut pending_idle: Option<IdleResource<R>>;
let mut state = self.state.lock();
state.active = state.active.saturating_sub(1);
if state.closed {
drop(state);
drop(resource);
continue;
}
let idle_since = (self.time_getter)();
pending_idle = Some(IdleResource {
resource,
idle_since,
created_at,
});
state.idle.reserve(1);
state.idle.push_back(
pending_idle
.take()
.expect("return payload prepared before idle insertion"),
);
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
if available > 0 && available - 1 < state.waiters.len() {
waiters_to_wake.push(state.waiters[available - 1].waker.clone_waker());
}
drop(state);
}
PoolReturn::Discard { hold_duration } => {
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_released(hold_duration);
metrics.record_destroyed(DestroyReason::Unhealthy);
}
let mut state = self.state.lock();
state.active = state.active.saturating_sub(1);
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
if available > 0 && available - 1 < state.waiters.len() {
waiters_to_wake.push(state.waiters[available - 1].waker.clone_waker());
}
drop(state);
}
}
}
for waker in waiters_to_wake {
waker.wake();
}
}
fn try_get_idle(&self, waiter_id: Option<u64>) -> Option<(R, Instant)> {
let now = (self.time_getter)();
loop {
let mut retired_idle = Vec::new();
let mut state = self.state.lock();
#[cfg(feature = "metrics")]
let mut idle_timeout_evictions = 0u64;
#[cfg(feature = "metrics")]
let mut max_lifetime_evictions = 0u64;
let idle_len = state.idle.len();
let mut expired_count = 0usize;
for idle in &state.idle {
let idle_ok =
now.saturating_duration_since(idle.idle_since) < self.config.idle_timeout;
let lifetime_ok =
now.saturating_duration_since(idle.created_at) < self.config.max_lifetime;
if !idle_ok || !lifetime_ok {
expired_count += 1;
#[cfg(feature = "metrics")]
if !idle_ok {
idle_timeout_evictions += 1;
} else {
max_lifetime_evictions += 1;
}
}
}
retired_idle.reserve(expired_count);
if expired_count > 0 {
for _ in 0..idle_len {
let idle = state
.idle
.pop_front()
.expect("idle partition length captured under state lock");
let idle_ok =
now.saturating_duration_since(idle.idle_since) < self.config.idle_timeout;
let lifetime_ok =
now.saturating_duration_since(idle.created_at) < self.config.max_lifetime;
if idle_ok && lifetime_ok {
state.idle.push_back(idle);
} else {
retired_idle.push(idle);
}
}
drop(state);
drop(retired_idle);
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
for _ in 0..idle_timeout_evictions {
metrics.record_destroyed(DestroyReason::IdleTimeout);
}
for _ in 0..max_lifetime_evictions {
metrics.record_destroyed(DestroyReason::MaxLifetime);
}
}
continue;
}
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
let pos = waiter_id.map_or_else(
|| state.waiters.len(),
|id| {
state
.waiters
.iter()
.position(|w| w.id == id)
.unwrap_or(state.waiters.len())
},
);
let result = if pos < available && !state.idle.is_empty() {
let active = state.active + 1;
let total_acquisitions = state.total_acquisitions + 1;
let idle = state
.idle
.pop_front()
.expect("non-empty idle queue checked under state lock");
state.active = active;
state.total_acquisitions = total_acquisitions;
Some((idle.resource, idle.created_at))
} else {
None
};
drop(state);
return result;
}
}
fn reserve_create_slot(&self, waiter_id: Option<u64>) -> Option<CreateSlotClaim> {
let mut state = self.state.lock();
let total = state.active + state.idle.len() + state.creating;
if state.closed || total >= self.config.max_size {
return None;
}
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
let pos = waiter_id.map_or_else(
|| state.waiters.len(),
|id| {
state
.waiters
.iter()
.position(|w| w.id == id)
.unwrap_or(state.waiters.len())
},
);
if pos >= available {
return None;
}
state.creating += 1;
let retired_waker = (waiter_id.is_some() && pos < state.waiters.len()).then(|| {
state
.waiters
.remove(pos)
.expect("waiter position was validated")
.waker
});
drop(state);
Some(CreateSlotClaim { retired_waker })
}
fn release_create_slot(&self) {
let waker = {
let mut state = self.state.lock();
state.creating = state.creating.saturating_sub(1);
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
if available > 0 && available - 1 < state.waiters.len() {
Some(state.waiters[available - 1].waker.clone_waker())
} else {
None
}
};
if let Some(waker) = waker {
waker.wake();
}
}
fn commit_create_slot(&self) -> bool {
let mut state = self.state.lock();
let creating = state.creating.saturating_sub(1);
let total_created = state.total_created + 1;
if state.closed {
state.creating = creating;
state.total_created = total_created;
return false;
}
let active = state.active + 1;
let total_acquisitions = state.total_acquisitions + 1;
state.creating = creating;
state.total_created = total_created;
state.active = active;
state.total_acquisitions = total_acquisitions;
true
}
fn commit_create_slot_as_idle(&self, resource: R) {
#[allow(clippy::needless_late_init)]
let mut pending_idle: Option<IdleResource<R>>;
let waker = {
let mut state = self.state.lock();
state.creating = state.creating.saturating_sub(1);
state.total_created += 1;
if state.closed {
drop(state);
return;
}
let now = (self.time_getter)();
pending_idle = Some(IdleResource {
resource,
idle_since: now,
created_at: now,
});
state.idle.reserve(1);
state.idle.push_back(
pending_idle
.take()
.expect("warmup payload prepared before idle insertion"),
);
let total = state.active + state.idle.len() + state.creating;
let available = state.idle.len() + self.config.max_size.saturating_sub(total);
if available > 0 && available - 1 < state.waiters.len() {
Some(state.waiters[available - 1].waker.clone_waker())
} else {
None
}
};
if let Some(waker) = waker {
waker.wake();
}
}
async fn create_resource(&self) -> Result<R, PoolError> {
let fut = self.factory.create();
fut.await.map_err(|e| PoolError::CreateFailed(e.into()))
}
fn remaining_acquire_timeout(
&self,
cx: &Cx,
acquire_start: crate::types::Time,
now: crate::types::Time,
) -> Result<Duration, PoolError> {
let elapsed = Duration::from_nanos(now.duration_since(acquire_start));
if elapsed >= self.config.acquire_timeout {
return Err(PoolError::Timeout);
}
let remaining = self.config.acquire_timeout.saturating_sub(elapsed);
if let Some(budget_remaining) = cx.budget().remaining_time(now) {
if budget_remaining.is_zero() {
return Err(PoolError::Cancelled);
}
Ok(remaining.min(budget_remaining))
} else {
Ok(remaining)
}
}
fn detach_waiter(&self, id: u64) -> Option<DeferredWaker> {
let mut state = self.state.lock();
state
.waiters
.iter()
.position(|waiter| waiter.id == id)
.and_then(|position| state.waiters.remove(position))
.map(|waiter| waiter.waker)
}
fn detach_return_waker(&self, id: u64) -> (Option<DeferredWaker>, Option<Waker>) {
let mut retired_waker = None;
let next_dispatcher = {
let mut wakers = self.return_wakers.lock();
if let Some(position) = wakers.iter().position(|(waiter_id, _)| *waiter_id == id) {
let was_dispatcher = position == 0;
retired_waker = Some(wakers.remove(position).1);
if was_dispatcher {
wakers.first().map(|(_, next)| next.clone_waker())
} else {
None
}
} else {
None
}
};
(retired_waker, next_dispatcher)
}
fn detach_waiter_registrations(&self, id: u64) -> DetachedPoolWakers {
let state_waker = self.detach_waiter(id);
let (return_waker, next_dispatcher) = self.detach_return_waker(id);
DetachedPoolWakers {
state_waker,
return_waker,
next_dispatcher,
}
}
#[cfg_attr(not(feature = "metrics"), allow(unused_variables))]
fn record_wait_time(&self, wait_duration: Duration) {
if wait_duration.is_zero() {
return;
}
let mut state = self.state.lock();
state.total_wait_time = state
.total_wait_time
.checked_add(wait_duration)
.unwrap_or(Duration::MAX);
drop(state);
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_wait(wait_duration);
}
}
#[cfg(feature = "metrics")]
fn update_metrics_gauges(&self) {
if let Some(ref metrics) = self.metrics {
let stats = {
let state = self.state.lock();
PoolStats {
active: state.active,
idle: state.idle.len(),
total: state.active + state.idle.len() + state.creating,
max_size: self.config.max_size,
waiters: state.waiters.len(),
total_acquisitions: state.total_acquisitions,
total_wait_time: state.total_wait_time,
}
};
metrics.update_gauges(&stats);
}
}
}
impl<R, F> Pool for GenericPool<R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
type Resource = R;
type Error = PoolError;
#[cfg_attr(not(feature = "metrics"), allow(unused_variables))]
#[allow(clippy::too_many_lines)]
fn acquire<'a>(
&'a self,
cx: &'a Cx,
) -> PoolFuture<'a, Result<PooledResource<Self::Resource>, Self::Error>> {
Box::pin(async move {
struct WaiterCleanup<'a, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
pool: &'a GenericPool<R, F>,
waiter_id: Option<u64>,
}
impl<R, F> Drop for WaiterCleanup<'_, R, F>
where
R: Send + 'static,
F: AsyncResourceFactory<Resource = R>,
{
fn drop(&mut self) {
if let Some(id) = self.waiter_id {
let mut retired_state_waker = None;
let marginal_waker = {
let mut state = self.pool.state.lock();
let pos = state.waiters.iter().position(|w| w.id == id);
if let Some(p) = pos {
retired_state_waker =
state.waiters.remove(p).map(|waiter| waiter.waker);
}
if state.closed {
None
} else {
let total_including_creating =
state.active + state.idle.len() + state.creating;
let available = state.idle.len()
+ self
.pool
.config
.max_size
.saturating_sub(total_including_creating);
pos.and_then(|p| {
if p < available
&& available > 0
&& available - 1 < state.waiters.len()
{
Some(state.waiters[available - 1].waker.clone_waker())
} else {
None
}
})
}
};
let (retired_return_waker, next_dispatcher) =
self.pool.detach_return_waker(id);
if let Some(w) = marginal_waker {
w.wake();
}
DetachedPoolWakers {
state_waker: retired_state_waker,
return_waker: retired_return_waker,
next_dispatcher,
}
.retire();
}
self.pool.process_returns();
}
}
let get_now = || {
cx.timer_driver()
.map_or_else(crate::time::wall_now, |d| d.now())
};
let acquire_start = get_now();
let mut cleanup = WaiterCleanup {
pool: self,
waiter_id: None,
};
loop {
self.process_returns();
if cx.checkpoint().is_err() {
return Err(PoolError::Cancelled);
}
if self.closed.load(Ordering::Acquire) {
return Err(PoolError::Closed);
}
while let Some((resource, created_at)) = self.try_get_idle(cleanup.waiter_id) {
let mut checkout = ActiveCheckoutGuard::new(self);
let is_healthy = if self.config.health_check_on_acquire {
checkout.begin_health_check();
let healthy = self.is_healthy(&resource);
checkout.finish_health_check();
healthy
} else {
true
};
if !is_healthy {
checkout.reject_unhealthy();
continue;
}
if let Some(id) = cleanup.waiter_id {
let detached = self.detach_waiter_registrations(id);
cleanup.waiter_id = None;
detached.retire();
}
let acquire_duration =
Duration::from_nanos(get_now().duration_since(acquire_start));
let acquired_at = (self.time_getter)();
let pooled = PooledResource::new_with_timestamps(
resource,
self.return_tx.clone(),
acquired_at,
created_at,
self.time_getter,
)
.with_return_notify(Arc::clone(&self.return_wakers));
checkout.commit();
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_acquired(acquire_duration);
self.update_metrics_gauges();
}
return Ok(pooled);
}
if let Some(create_slot) =
CreateSlotReservation::try_reserve(self, cleanup.waiter_id)
{
if let Some(id) = cleanup.waiter_id {
let detached = self.detach_waiter_registrations(id);
cleanup.waiter_id = None;
detached.retire();
}
let now = get_now();
let remaining = match self.remaining_acquire_timeout(cx, acquire_start, now) {
Ok(remaining) => remaining,
Err(PoolError::Timeout) => {
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_timeout(Duration::from_nanos(
now.duration_since(acquire_start),
));
}
return Err(PoolError::Timeout);
}
Err(PoolError::Cancelled) => return Err(PoolError::Cancelled),
Err(other) => return Err(other),
};
if remaining.is_zero() {
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_timeout(Duration::from_nanos(
now.duration_since(acquire_start),
));
}
return if cx.checkpoint().is_err() {
Err(PoolError::Cancelled)
} else {
Err(PoolError::Timeout)
};
}
let create_result =
crate::time::timeout(now, remaining, self.create_resource()).await;
let resource = match create_result {
Ok(Ok(res)) => res,
Ok(Err(e)) => return Err(e),
Err(_) => {
if cx.checkpoint().is_err() {
return Err(PoolError::Cancelled);
}
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_timeout(Duration::from_nanos(
get_now().duration_since(acquire_start),
));
}
return Err(PoolError::Timeout);
}
};
let committed = create_slot.commit();
let checkout = committed.then(|| ActiveCheckoutGuard::new(self));
let acquire_duration =
Duration::from_nanos(get_now().duration_since(acquire_start));
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_created();
}
if !committed {
#[cfg(feature = "metrics")]
self.update_metrics_gauges();
return Err(PoolError::Closed);
}
let acquired_at = (self.time_getter)();
let pooled = PooledResource::new_with_timestamps(
resource,
self.return_tx.clone(),
acquired_at,
acquired_at,
self.time_getter,
)
.with_return_notify(Arc::clone(&self.return_wakers));
checkout
.expect("committed creation slot arms active checkout guard")
.commit();
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_acquired(acquire_duration);
self.update_metrics_gauges();
}
return Ok(pooled);
}
let now = get_now();
let remaining = match self.remaining_acquire_timeout(cx, acquire_start, now) {
Ok(remaining) => remaining,
Err(PoolError::Timeout) => {
let elapsed = Duration::from_nanos(now.duration_since(acquire_start));
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_timeout(elapsed);
}
return Err(PoolError::Timeout);
}
Err(PoolError::Cancelled) => return Err(PoolError::Cancelled),
Err(other) => return Err(other),
};
if let Err(_e) = cx.checkpoint() {
return Err(PoolError::Cancelled);
}
let wait_started = now;
let wait_fut = WaitForNotification {
pool: self,
waiter_id: &mut cleanup.waiter_id,
cx,
completed: false,
};
if crate::time::timeout(now, remaining, wait_fut)
.await
.is_err()
{
let wait_duration =
Duration::from_nanos(get_now().duration_since(wait_started));
if cx.checkpoint().is_err() {
self.record_wait_time(wait_duration);
return Err(PoolError::Cancelled);
}
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.record_timeout(wait_duration);
}
self.record_wait_time(wait_duration);
return Err(PoolError::Timeout);
}
self.record_wait_time(Duration::from_nanos(get_now().duration_since(wait_started)));
}
})
}
#[cfg_attr(not(feature = "metrics"), allow(unused_variables))]
fn try_acquire(&self) -> Option<PooledResource<Self::Resource>> {
let acquire_start = (self.time_getter)();
self.process_returns();
if self.closed.load(Ordering::Acquire) {
return None;
}
while let Some((resource, created_at)) = self.try_get_idle(None) {
let mut checkout = ActiveCheckoutGuard::new(self);
let is_healthy = if self.config.health_check_on_acquire {
checkout.begin_health_check();
let healthy = self.is_healthy(&resource);
checkout.finish_health_check();
healthy
} else {
true
};
if !is_healthy {
checkout.reject_unhealthy();
continue;
}
#[cfg(feature = "metrics")]
let acquire_duration = self
.metrics
.as_ref()
.map(|_| (self.time_getter)().saturating_duration_since(acquire_start));
let acquired_at = (self.time_getter)();
let pooled = PooledResource::new_with_timestamps(
resource,
self.return_tx.clone(),
acquired_at,
created_at,
self.time_getter,
)
.with_return_notify(Arc::clone(&self.return_wakers));
checkout.commit();
#[cfg(feature = "metrics")]
if let (Some(metrics), Some(acquire_duration)) = (&self.metrics, acquire_duration) {
metrics.record_acquired(acquire_duration);
self.update_metrics_gauges();
}
return Some(pooled);
}
None
}
fn stats(&self) -> PoolStats {
self.process_returns();
let pool_stats = {
let state = self.state.lock();
PoolStats {
active: state.active,
idle: state.idle.len(),
total: state.active + state.idle.len() + state.creating,
max_size: self.config.max_size,
waiters: state.waiters.len(),
total_acquisitions: state.total_acquisitions,
total_wait_time: state.total_wait_time,
}
};
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
metrics.update_gauges(&pool_stats);
}
pool_stats
}
fn close(&self) -> PoolFuture<'_, ()> {
Box::pin(async move {
let (waiters, idle_resources) = {
let mut state = self.state.lock();
state.closed = true;
self.closed.store(true, Ordering::Release);
let waiters = std::mem::take(&mut state.waiters);
let idle_resources = std::mem::take(&mut state.idle);
(waiters, idle_resources)
};
#[cfg(feature = "metrics")]
let idle_count = idle_resources.len();
for waiter in waiters {
waiter.waker.clone_waker().wake();
}
#[cfg(feature = "metrics")]
if let Some(ref metrics) = self.metrics {
for _ in 0..idle_count {
metrics.record_destroyed(DestroyReason::Unhealthy);
}
self.update_metrics_gauges();
}
drop(idle_resources);
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DestroyReason {
Unhealthy,
IdleTimeout,
MaxLifetime,
}
impl DestroyReason {
#[must_use]
pub const fn as_label(&self) -> &'static str {
match self {
Self::Unhealthy => "unhealthy",
Self::IdleTimeout => "idle_timeout",
Self::MaxLifetime => "max_lifetime",
}
}
}
#[cfg(feature = "metrics")]
mod pool_metrics {
use super::{DestroyReason, Duration, PoolStats};
use opentelemetry::KeyValue;
use opentelemetry::metrics::{Counter, Histogram, Meter, ObservableGauge};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Default)]
pub struct PoolMetricsState {
pub size: AtomicU64,
pub active: AtomicU64,
pub idle: AtomicU64,
pub pending: AtomicU64,
}
impl PoolMetricsState {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn update_from_stats(&self, stats: &PoolStats) {
self.size.store(stats.total as u64, Ordering::Relaxed);
self.active.store(stats.active as u64, Ordering::Relaxed);
self.idle.store(stats.idle as u64, Ordering::Relaxed);
self.pending.store(stats.waiters as u64, Ordering::Relaxed);
}
}
#[derive(Clone)]
pub struct PoolMetrics {
#[allow(dead_code)]
size: ObservableGauge<u64>,
#[allow(dead_code)]
active: ObservableGauge<u64>,
#[allow(dead_code)]
idle: ObservableGauge<u64>,
#[allow(dead_code)]
pending: ObservableGauge<u64>,
acquired_total: Counter<u64>,
released_total: Counter<u64>,
created_total: Counter<u64>,
destroyed_total: Counter<u64>,
timeouts_total: Counter<u64>,
acquire_duration: Histogram<f64>,
hold_duration: Histogram<f64>,
wait_duration: Histogram<f64>,
state: Arc<PoolMetricsState>,
}
impl std::fmt::Debug for PoolMetrics {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PoolMetrics")
.field("state", &self.state)
.finish_non_exhaustive()
}
}
impl PoolMetrics {
#[must_use]
pub fn new(meter: &Meter) -> Self {
let state = Arc::new(PoolMetricsState::new());
let size = meter
.u64_observable_gauge("asupersync.pool.size")
.with_description("Current pool size (active + idle)")
.with_callback({
let state = Arc::clone(&state);
move |observer| {
observer.observe(state.size.load(Ordering::Relaxed), &[]);
}
})
.build();
let active = meter
.u64_observable_gauge("asupersync.pool.active")
.with_description("Currently checked-out resources")
.with_callback({
let state = Arc::clone(&state);
move |observer| {
observer.observe(state.active.load(Ordering::Relaxed), &[]);
}
})
.build();
let idle = meter
.u64_observable_gauge("asupersync.pool.idle")
.with_description("Available idle resources")
.with_callback({
let state = Arc::clone(&state);
move |observer| {
observer.observe(state.idle.load(Ordering::Relaxed), &[]);
}
})
.build();
let pending = meter
.u64_observable_gauge("asupersync.pool.pending")
.with_description("Waiters in queue")
.with_callback({
let state = Arc::clone(&state);
move |observer| {
observer.observe(state.pending.load(Ordering::Relaxed), &[]);
}
})
.build();
let acquired_total = meter
.u64_counter("asupersync.pool.acquired_total")
.with_description("Total successful acquires")
.build();
let released_total = meter
.u64_counter("asupersync.pool.released_total")
.with_description("Total returns to pool")
.build();
let created_total = meter
.u64_counter("asupersync.pool.created_total")
.with_description("Resources created")
.build();
let destroyed_total = meter
.u64_counter("asupersync.pool.destroyed_total")
.with_description("Resources destroyed")
.build();
let timeouts_total = meter
.u64_counter("asupersync.pool.timeouts_total")
.with_description("Acquire timeouts")
.build();
let acquire_duration = meter
.f64_histogram("asupersync.pool.acquire_duration_seconds")
.with_description("Time to acquire a resource")
.build();
let hold_duration = meter
.f64_histogram("asupersync.pool.hold_duration_seconds")
.with_description("Time resource is held")
.build();
let wait_duration = meter
.f64_histogram("asupersync.pool.wait_duration_seconds")
.with_description("Time waiting in queue")
.build();
Self {
size,
active,
idle,
pending,
acquired_total,
released_total,
created_total,
destroyed_total,
timeouts_total,
acquire_duration,
hold_duration,
wait_duration,
state,
}
}
#[must_use]
pub fn state(&self) -> &Arc<PoolMetricsState> {
&self.state
}
pub fn record_acquired(&self, pool_name: &str, duration: Duration) {
let labels = [KeyValue::new("pool_name", pool_name.to_string())];
self.acquired_total.add(1, &labels);
self.acquire_duration
.record(duration.as_secs_f64(), &labels);
}
pub fn record_released(&self, pool_name: &str, hold_duration: Duration) {
let labels = [KeyValue::new("pool_name", pool_name.to_string())];
self.released_total.add(1, &labels);
self.hold_duration
.record(hold_duration.as_secs_f64(), &labels);
}
pub fn record_created(&self, pool_name: &str) {
let labels = [KeyValue::new("pool_name", pool_name.to_string())];
self.created_total.add(1, &labels);
}
pub fn record_destroyed(&self, pool_name: &str, reason: DestroyReason) {
let labels = [
KeyValue::new("pool_name", pool_name.to_string()),
KeyValue::new("reason", reason.as_label()),
];
self.destroyed_total.add(1, &labels);
}
pub fn record_timeout(&self, pool_name: &str, wait_duration: Duration) {
let labels = [KeyValue::new("pool_name", pool_name.to_string())];
self.timeouts_total.add(1, &labels);
self.wait_duration
.record(wait_duration.as_secs_f64(), &labels);
}
pub fn record_wait(&self, pool_name: &str, wait_duration: Duration) {
let labels = [KeyValue::new("pool_name", pool_name.to_string())];
self.wait_duration
.record(wait_duration.as_secs_f64(), &labels);
}
pub fn update_gauges(&self, stats: &PoolStats) {
self.state.update_from_stats(stats);
}
#[must_use]
pub fn handle(&self, pool_name: impl Into<String>) -> PoolMetricsHandle {
let pool_name = pool_name.into();
let labels = [KeyValue::new("pool_name", pool_name.clone())];
PoolMetricsHandle {
metrics: self.clone(),
pool_name,
labels,
}
}
}
#[derive(Clone)]
pub struct PoolMetricsHandle {
metrics: PoolMetrics,
pool_name: String,
labels: [KeyValue; 1],
}
impl std::fmt::Debug for PoolMetricsHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PoolMetricsHandle")
.field("pool_name", &self.pool_name)
.finish_non_exhaustive()
}
}
impl PoolMetricsHandle {
#[must_use]
pub fn pool_name(&self) -> &str {
&self.pool_name
}
pub fn record_acquired(&self, duration: Duration) {
self.metrics.acquired_total.add(1, &self.labels);
self.metrics
.acquire_duration
.record(duration.as_secs_f64(), &self.labels);
}
pub fn record_released(&self, hold_duration: Duration) {
self.metrics.released_total.add(1, &self.labels);
self.metrics
.hold_duration
.record(hold_duration.as_secs_f64(), &self.labels);
}
pub fn record_created(&self) {
self.metrics.created_total.add(1, &self.labels);
}
pub fn record_destroyed(&self, reason: DestroyReason) {
let labels = [
self.labels[0].clone(),
KeyValue::new("reason", reason.as_label()),
];
self.metrics.destroyed_total.add(1, &labels);
}
pub fn record_timeout(&self, wait_duration: Duration) {
self.metrics.timeouts_total.add(1, &self.labels);
self.metrics
.wait_duration
.record(wait_duration.as_secs_f64(), &self.labels);
}
pub fn record_wait(&self, wait_duration: Duration) {
self.metrics
.wait_duration
.record(wait_duration.as_secs_f64(), &self.labels);
}
pub fn update_gauges(&self, stats: &PoolStats) {
self.metrics.update_gauges(stats);
}
#[must_use]
pub fn state(&self) -> &Arc<PoolMetricsState> {
self.metrics.state()
}
}
}
#[cfg(feature = "metrics")]
pub use pool_metrics::{PoolMetrics, PoolMetricsHandle, PoolMetricsState};
#[cfg(test)]
include!("pool_tests.rs");