use std::ops::Deref;
use std::ops::DerefMut;
use std::sync::Arc;
use std::sync::Weak;
use crate::internal::mutex::Mutex;
use crate::pool::ManageObject;
use crate::pool::ObjectStatus;
use crate::pool::QueueStrategy;
use crate::pool::RecycleCancelledStrategy;
use crate::pool::RetainResult;
use crate::pool::state::ObjectState;
use crate::pool::state::PoolState;
use crate::semaphore::OwnedSemaphorePermit;
use crate::semaphore::Semaphore;
#[derive(Clone, Copy, Debug)]
#[non_exhaustive]
pub struct PoolConfig {
pub max_size: usize,
pub queue_strategy: QueueStrategy,
pub recycle_cancelled_strategy: RecycleCancelledStrategy,
}
impl PoolConfig {
pub fn new(max_size: usize) -> Self {
Self {
max_size,
queue_strategy: QueueStrategy::default(),
recycle_cancelled_strategy: RecycleCancelledStrategy::default(),
}
}
#[must_use = "this method returns the updated pool configuration"]
pub fn with_queue_strategy(mut self, queue_strategy: QueueStrategy) -> Self {
self.queue_strategy = queue_strategy;
self
}
#[must_use = "this method returns the updated pool configuration"]
pub fn with_recycle_cancelled_strategy(
mut self,
recycle_cancelled_strategy: RecycleCancelledStrategy,
) -> Self {
self.recycle_cancelled_strategy = recycle_cancelled_strategy;
self
}
}
#[derive(Clone, Copy, Debug)]
#[non_exhaustive]
pub struct PoolStatus {
pub max_size: usize,
pub current_size: usize,
pub idle_count: usize,
}
pub struct Pool<M: ManageObject> {
config: PoolConfig,
manager: M,
permits: Arc<Semaphore>,
slots: Mutex<PoolState<M::Object>>,
}
impl<M> std::fmt::Debug for Pool<M>
where
M: ManageObject,
M::Object: std::fmt::Debug,
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Pool")
.field("slots", &self.slots)
.field("config", &self.config)
.field("permits", &self.permits)
.finish()
}
}
impl<M: ManageObject> Pool<M> {
pub fn new(config: PoolConfig, manager: M) -> Arc<Self> {
assert!(
config.max_size > 0,
"bounded pool max_size must be greater than zero"
);
let permits = Arc::new(Semaphore::new(config.max_size));
let slots = Mutex::new(PoolState::new());
Arc::new(Self {
config,
manager,
permits,
slots,
})
}
pub async fn replenish_to(&self, target_idle: usize) -> Result<usize, M::Error> {
let target_idle = target_idle.min(self.config.max_size);
let Some(mut reservation) = ReplenishReservation::reserve_up_to(&self.permits, target_idle)
else {
return Ok(0);
};
let (idle_count, available_slots) = {
let slots = self.slots.lock();
let idle_count = slots.idle_count();
let uncommitted_capacity = self
.permits
.available_permits()
.checked_add(reservation.permits())
.expect("invariant broken: semaphore capacity must not overflow");
let available_slots = uncommitted_capacity.saturating_sub(idle_count);
(idle_count, available_slots)
};
let to_create = target_idle
.saturating_sub(idle_count)
.min(reservation.permits())
.min(available_slots);
reservation.release(reservation.permits() - to_create);
let mut replenished = 0;
for _ in 0..to_create {
let object = self.manager.create().await?;
{
let mut slots = self.slots.lock();
slots.add_idle(ObjectState::new(object));
}
replenished += 1;
reservation.release(1);
}
Ok(replenished)
}
pub async fn get(self: &Arc<Self>) -> Result<Object<M>, M::Error> {
let permit = self.permits.clone().acquire_owned(1).await;
let object = loop {
let existing = self.slots.lock().pop(self.config.queue_strategy);
match existing {
None => {
let object = self.manager.create().await?;
let state = ObjectState::new(object);
self.slots.lock().add_active();
break Object {
state: Some(state),
permit,
pool: Arc::downgrade(self),
};
}
Some(object) => {
let mut unready_object = UnreadyObject {
state: Some(object),
pool: Arc::downgrade(self),
recycle_cancelled_strategy: self.config.recycle_cancelled_strategy,
};
let state = unready_object.state();
let status = state.status;
if self
.manager
.is_recyclable(&mut state.o, &status)
.await
.is_ok()
{
state.status.mark_recycled();
break unready_object.ready(permit);
} else {
unready_object.detach();
}
}
};
};
Ok(object)
}
pub fn retain(
&self,
f: impl FnMut(&mut M::Object, ObjectStatus) -> bool,
) -> RetainResult<M::Object> {
let mut result = {
let mut slots = self.slots.lock();
slots.retain(f)
};
for object in &mut result.removed {
self.manager.on_detached(object);
}
result
}
pub fn status(&self) -> PoolStatus {
let slots = self.slots.lock();
PoolStatus {
max_size: self.config.max_size,
current_size: slots.current_size(),
idle_count: slots.idle_count(),
}
}
fn return_object(&self, mut state: ObjectState<M::Object>) {
state.status.mark_returned();
self.restore_idle(state);
}
fn restore_idle(&self, state: ObjectState<M::Object>) {
let mut slots = self.slots.lock();
assert!(
slots.current_size() <= self.config.max_size,
"invariant broken: current_size <= max_size (actual: {} <= {})",
slots.current_size(),
self.config.max_size,
);
slots.return_idle(state);
}
fn detach_object(&self, o: &mut M::Object) {
let mut slots = self.slots.lock();
assert!(
slots.current_size() <= self.config.max_size,
"invariant broken: current_size <= max_size (actual: {} <= {})",
slots.current_size(),
self.config.max_size,
);
slots.detach();
drop(slots);
self.manager.on_detached(o);
}
}
struct ReplenishReservation<'a> {
semaphore: &'a Semaphore,
permits: usize,
}
impl<'a> ReplenishReservation<'a> {
fn reserve_up_to(semaphore: &'a Semaphore, up_to: usize) -> Option<Self> {
let permits = semaphore.drain_permits(up_to);
(permits != 0).then_some(Self { semaphore, permits })
}
fn permits(&self) -> usize {
self.permits
}
fn release(&mut self, permits: usize) {
assert!(
permits <= self.permits,
"cannot release more permits than this reservation holds"
);
self.permits -= permits;
self.semaphore.release(permits);
}
}
impl Drop for ReplenishReservation<'_> {
fn drop(&mut self) {
self.semaphore.release(self.permits);
}
}
pub struct Object<M: ManageObject> {
state: Option<ObjectState<M::Object>>,
permit: OwnedSemaphorePermit,
pool: Weak<Pool<M>>,
}
impl<M> std::fmt::Debug for Object<M>
where
M: ManageObject,
M::Object: std::fmt::Debug,
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Object")
.field("state", &self.state)
.field("permit", &self.permit)
.finish()
}
}
impl<M: ManageObject> Drop for Object<M> {
fn drop(&mut self) {
if let Some(state) = self.state.take() {
if let Some(pool) = self.pool.upgrade() {
pool.return_object(state);
}
}
}
}
impl<M: ManageObject> Deref for Object<M> {
type Target = M::Object;
fn deref(&self) -> &M::Object {
&self.state.as_ref().unwrap().o
}
}
impl<M: ManageObject> DerefMut for Object<M> {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.state.as_mut().unwrap().o
}
}
impl<M: ManageObject> AsRef<M::Object> for Object<M> {
fn as_ref(&self) -> &M::Object {
self
}
}
impl<M: ManageObject> AsMut<M::Object> for Object<M> {
fn as_mut(&mut self) -> &mut M::Object {
self
}
}
impl<M: ManageObject> Object<M> {
pub fn detach(mut self) -> M::Object {
let mut o = self.state.take().unwrap().o;
if let Some(pool) = self.pool.upgrade() {
pool.detach_object(&mut o);
}
o
}
pub fn status(&self) -> ObjectStatus {
self.state.as_ref().unwrap().status
}
}
struct UnreadyObject<M: ManageObject> {
state: Option<ObjectState<M::Object>>,
pool: Weak<Pool<M>>,
recycle_cancelled_strategy: RecycleCancelledStrategy,
}
impl<M: ManageObject> Drop for UnreadyObject<M> {
fn drop(&mut self) {
if let Some(mut state) = self.state.take() {
if let Some(pool) = self.pool.upgrade() {
match self.recycle_cancelled_strategy {
RecycleCancelledStrategy::Detach => {
pool.detach_object(&mut state.o);
}
RecycleCancelledStrategy::ReturnToPool => {
pool.restore_idle(state);
}
}
}
}
}
}
impl<M: ManageObject> UnreadyObject<M> {
fn ready(mut self, permit: OwnedSemaphorePermit) -> Object<M> {
let state = Some(self.state.take().unwrap());
let pool = self.pool.clone();
Object {
state,
permit,
pool,
}
}
fn detach(&mut self) {
if let Some(mut state) = self.state.take() {
if let Some(pool) = self.pool.upgrade() {
pool.detach_object(&mut state.o);
}
}
}
fn state(&mut self) -> &mut ObjectState<M::Object> {
self.state.as_mut().unwrap()
}
}