use super::*;
impl DbPool {
#[cfg(feature = "metrics")]
#[inline]
pub(super) fn record_acquire_duration(&self, start: Instant) {
if let Some(collector) = self
.inner
.metrics_collector
.read()
.expect("metrics_collector lock")
.clone()
{
collector.record_connection_acquire_duration(start.elapsed());
}
}
#[cfg(not(feature = "metrics"))]
#[inline]
pub(super) fn record_acquire_duration(&self, _start: Instant) {}
#[cfg(feature = "metrics")]
#[inline]
pub(super) fn record_acquire_timeout(&self, start: Instant) {
let elapsed_ms = start.elapsed().as_millis() as u64;
if let Some(collector) = self
.inner
.metrics_collector
.read()
.expect("metrics_collector lock")
.clone()
{
collector.record_connection_timeout_level(elapsed_ms);
}
}
#[cfg(not(feature = "metrics"))]
#[inline]
pub(super) fn record_acquire_timeout(&self, _start: Instant) {}
pub(super) fn update_max_waiters(&self, current_waiters: u32) {
let mut current = self.inner.max_waiters.load(Ordering::Acquire);
while current_waiters > current {
match self.inner.max_waiters.compare_exchange(
current,
current_waiters,
Ordering::SeqCst,
Ordering::Acquire,
) {
Ok(_) => return,
Err(observed) => {
current = observed;
}
}
}
}
#[cfg(any(feature = "auto-migrate", all(feature = "postgres", feature = "copy")))]
pub(crate) fn release_connection(&self, conn: DbConnection) {
DbPoolInner::release_connection(&self.inner, conn);
}
pub fn status(&self) -> PoolStatus {
let total = self.inner.total_count.load(Ordering::SeqCst);
let active = self.inner.active_count.load(Ordering::SeqCst);
let wait_count = self.inner.wait_count.load(Ordering::SeqCst);
let max_waiters = self.inner.max_waiters.load(Ordering::SeqCst);
let borrow_count = self.inner.borrow_count.load(Ordering::SeqCst);
let max_active = self.inner.max_active.load(Ordering::SeqCst);
PoolStatus {
total,
active,
idle: total.saturating_sub(active),
wait_count,
max_waiters,
borrow_count,
max_active,
}
}
#[cfg(feature = "metrics")]
pub fn pool_metrics(&self) -> PoolMetrics {
let wait_count = self.inner.wait_count.load(Ordering::SeqCst);
let max_waiters = self.inner.max_waiters.load(Ordering::SeqCst);
let collector = self
.inner
.metrics_collector
.read()
.expect("metrics_collector lock")
.clone();
if let Some(collector) = collector {
let stats = collector.connection_acquire_stats();
PoolMetrics {
slow_acquires: stats.slow_acquires,
timeout_errors: stats.timeout_warn + stats.timeout_error + stats.timeout_critical,
critical_timeouts: stats.timeout_critical,
wait_count,
max_waiters,
}
} else {
PoolMetrics {
slow_acquires: 0,
timeout_errors: 0,
critical_timeouts: 0,
wait_count,
max_waiters,
}
}
}
pub fn config(&self) -> &DbConfig {
&self.inner.config
}
#[cfg(feature = "auto-migrate")]
pub async fn run_auto_migrate(&self) -> Result<u32, DbError> {
if let Some(ref migrations_dir) = self.inner.config.migrations_dir {
self.run_migrations(migrations_dir).await
} else {
Ok(0)
}
}
#[cfg(feature = "auto-migrate")]
pub async fn run_migrations(&self, migrations_dir: &std::path::Path) -> Result<u32, DbError> {
use crate::database::MigrationExecutor;
let db_type = self
.inner
.config
.database_type()
.map_err(|e| DbError::Config(e.to_string()))?;
let connection = self.acquire_connection().await?;
let outcome: Result<u32, DbError> = async {
let connection_for_migration = connection.as_sea_orm()?.clone();
let mut executor = MigrationExecutor::new(connection_for_migration, db_type);
executor.run_migrations(migrations_dir).await
}
.await;
self.release_connection(connection);
outcome
}
}