mod contents;
mod manager;
use std::sync::mpsc;
use std::thread;
use crate::config::PoolConfig;
use crate::connection::Connection;
use crate::error::Error;
use contents::PoolAcquireResponse;
use contents::PoolContents;
pub(crate) use contents::PoolContentsRef;
use manager::PoolManager;
use manager::PoolManagerRequest;
pub struct Pool {
contents_ref: PoolContentsRef,
bg_task: Option<thread::JoinHandle<()>>,
}
impl Pool {
fn check_open(&self) -> Result<(), Error> {
if self.bg_task.is_some() {
Ok(())
} else {
Err(Error::pool_not_open())
}
}
pub(crate) fn create(config: PoolConfig) -> Result<Self, Error> {
config.validate()?;
let mut actual_config = config;
if actual_config.cclass().is_none() {
let cclass = format!("RSO:{}", uuid::Uuid::new_v4());
actual_config = actual_config.set_cclass(cclass);
}
let manager_config = actual_config.clone();
let (tx, rx) = mpsc::channel();
let contents = PoolContents::new(actual_config, tx);
let contents_ref =
std::sync::Arc::new(std::sync::Mutex::new(contents));
let manager_contents_ref = contents_ref.clone();
let bg_task = thread::spawn(move || {
let mut manager =
PoolManager::new(manager_contents_ref, rx, manager_config);
manager.run();
});
Ok(Self {
contents_ref,
bg_task: Some(bg_task),
})
}
pub fn acquire(&self) -> Result<Connection, Error> {
self.check_open()?;
let resp = self.contents_ref.lock().unwrap().acquire();
let conn_impl_result = match resp {
PoolAcquireResponse::Connection(result) => result,
PoolAcquireResponse::Wait(channel) => channel.recv().unwrap(),
};
conn_impl_result.map(|conn_impl| {
Connection::create_pooled(conn_impl, &self.contents_ref)
})
}
pub fn busy_count(&self) -> Result<usize, Error> {
self.check_open()?;
Ok(self.contents_ref.lock().unwrap().busy_count())
}
pub fn close(&mut self) -> Result<(), Error> {
self.check_open()?;
self.contents_ref.lock().unwrap().close()?;
if let Some(handle) = self.bg_task.take() {
let _ = handle.join();
}
Ok(())
}
pub fn open_count(&self) -> Result<usize, Error> {
self.check_open()?;
Ok(self.contents_ref.lock().unwrap().open_count())
}
}
impl Drop for Pool {
fn drop(&mut self) {
let _ = self.close();
}
}