use do_memory_core::{Error, Result};
use libsql::Connection;
use parking_lot::RwLock;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, AtomicUsize};
use std::time::Duration;
use tokio::sync::OwnedSemaphorePermit;
#[derive(Debug, Clone)]
pub struct PoolConfig {
pub min_connections: usize,
pub max_connections: usize,
pub connection_timeout: Duration,
pub enable_health_check: bool,
pub health_check_timeout: Duration,
pub acquire_timeout_ms: u64,
pub idle_timeout_ms: u64,
}
impl Default for PoolConfig {
fn default() -> Self {
Self {
min_connections: 1,
max_connections: 10,
connection_timeout: Duration::from_secs(5),
enable_health_check: true,
health_check_timeout: Duration::from_secs(2),
acquire_timeout_ms: 5000,
idle_timeout_ms: 0,
}
}
}
#[derive(Debug)]
pub struct PoolMetrics {
pub active_connections: AtomicUsize,
pub total_acquired: AtomicU64,
pub total_wait_ms: AtomicU64,
pub reconnect_count: AtomicU64,
}
impl Default for PoolMetrics {
fn default() -> Self {
Self {
active_connections: AtomicUsize::new(0),
total_acquired: AtomicU64::new(0),
total_wait_ms: AtomicU64::new(0),
reconnect_count: AtomicU64::new(0),
}
}
}
impl PoolMetrics {
pub fn record_acquire(&self, wait_ms: u64) {
self.active_connections
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.total_acquired
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.total_wait_ms
.fetch_add(wait_ms, std::sync::atomic::Ordering::Relaxed);
}
pub fn record_release(&self) {
self.active_connections
.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
}
pub fn record_reconnect(&self) {
self.reconnect_count
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
pub fn snapshot(&self) -> PoolMetricsSnapshot {
PoolMetricsSnapshot {
active_connections: self
.active_connections
.load(std::sync::atomic::Ordering::Relaxed),
total_acquired: self
.total_acquired
.load(std::sync::atomic::Ordering::Relaxed),
total_wait_ms: self
.total_wait_ms
.load(std::sync::atomic::Ordering::Relaxed),
reconnect_count: self
.reconnect_count
.load(std::sync::atomic::Ordering::Relaxed),
}
}
}
#[derive(Debug, Clone, Default)]
pub struct PoolMetricsSnapshot {
pub active_connections: usize,
pub total_acquired: u64,
pub total_wait_ms: u64,
pub reconnect_count: u64,
}
#[derive(Debug, Clone, Default)]
pub struct PoolStatistics {
pub total_created: usize,
pub total_health_checks_passed: usize,
pub total_health_checks_failed: usize,
pub active_connections: usize,
pub total_wait_time_ms: u64,
pub total_checkouts: usize,
pub avg_wait_time_ms: u64,
}
impl PoolStatistics {
pub fn update_averages(&mut self) {
if self.total_checkouts > 0 {
self.avg_wait_time_ms = self.total_wait_time_ms / self.total_checkouts as u64;
}
}
}
#[derive(Debug)]
pub struct PooledConnection {
pub(super) connection: Option<Connection>,
pub(super) _permit: OwnedSemaphorePermit,
pub(super) stats: Arc<RwLock<PoolStatistics>>,
pub(super) metrics: Option<Arc<PoolMetrics>>,
}
impl PooledConnection {
pub fn connection(&self) -> Option<&Connection> {
self.connection.as_ref()
}
pub fn into_inner(mut self) -> Result<Connection> {
self.connection
.take()
.ok_or_else(|| Error::Storage("Connection already taken".to_string()))
}
}
impl Drop for PooledConnection {
fn drop(&mut self) {
let mut stats = self.stats.write();
if stats.active_connections > 0 {
stats.active_connections -= 1;
}
if let Some(ref metrics) = self.metrics {
metrics.record_release();
}
}
}