mirror of
https://github.com/bitechdev/ResolveSpec.git
synced 2026-10-01 04:21:58 +00:00
feat(dbmanager): report pool max size and total connections opened
Add MaxOpenConnections to provider and connection stats, exposed as the max state of dbmanager_connection_pool_size. Count physical dials through a counting driver connector and export dbmanager_connections_opened_total; existing_db pools derive an approximate count from sql.DBStats.
This commit is contained in:
@@ -319,9 +319,10 @@ fmt.Printf("Healthy: %d, Unhealthy: %d\n", stats.HealthyCount, stats.UnhealthyCo
|
||||
|
||||
// Per-connection stats
|
||||
for name, connStats := range stats.ConnectionStats {
|
||||
fmt.Printf("%s: %d open, %d in use, %d idle\n",
|
||||
fmt.Printf("%s: %d open (max %d, 0 = unlimited), %d in use, %d idle\n",
|
||||
name,
|
||||
connStats.OpenConnections,
|
||||
connStats.MaxOpenConnections,
|
||||
connStats.InUse,
|
||||
connStats.Idle)
|
||||
}
|
||||
@@ -340,7 +341,8 @@ The package automatically exports Prometheus metrics:
|
||||
|
||||
- `dbmanager_connections_total` - Total configured connections by type
|
||||
- `dbmanager_connection_status` - Connection health status (1=healthy, 0=unhealthy)
|
||||
- `dbmanager_connection_pool_size` - Connection pool statistics by state
|
||||
- `dbmanager_connections_opened_total` - Physical connections ever opened (counter; for a pool wrapped via `existing_db` this is derived from open + closed counts and can undercount)
|
||||
- `dbmanager_connection_pool_size` - Connection pool statistics by state (`open`, `idle`, `in_use`, `max`; `max` 0 = unlimited)
|
||||
- `dbmanager_connection_wait_count` - Times connections waited for availability
|
||||
- `dbmanager_connection_wait_duration_seconds` - Total wait duration
|
||||
- `dbmanager_health_check_duration_seconds` - Health check execution time
|
||||
|
||||
@@ -5,6 +5,8 @@ import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
)
|
||||
|
||||
func TestPostgresDSNEscapesCredentials(t *testing.T) {
|
||||
@@ -83,6 +85,18 @@ func TestSQLiteMemoryPoolPinned(t *testing.T) {
|
||||
if got := db.Stats().MaxOpenConnections; got != 1 {
|
||||
t.Fatalf("MaxOpenConnections = %d, want 1 for :memory:", got)
|
||||
}
|
||||
if got := conn.Stats().MaxOpenConnections; got != 1 {
|
||||
t.Fatalf("ConnectionStats.MaxOpenConnections = %d, want 1", got)
|
||||
}
|
||||
mgr.(*connectionManager).PublishMetrics()
|
||||
name := conn.Name()
|
||||
var m dto.Metric
|
||||
if err := connectionPoolSize.WithLabelValues(name, string(conn.Stats().Type), "max").Write(&m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := m.GetGauge().GetValue(); got != 1 {
|
||||
t.Fatalf("pool_size{state=max} = %v, want 1", got)
|
||||
}
|
||||
if _, err := db.Exec("CREATE TABLE t(a int)"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@@ -55,13 +55,15 @@ type ConnectionStats struct {
|
||||
HealthCheckStatus string
|
||||
|
||||
// SQL connection pool stats
|
||||
OpenConnections int
|
||||
InUse int
|
||||
Idle int
|
||||
WaitCount int64
|
||||
WaitDuration time.Duration
|
||||
MaxIdleClosed int64
|
||||
MaxLifetimeClosed int64
|
||||
OpenConnections int
|
||||
MaxOpenConnections int
|
||||
TotalOpened int64 // physical connections ever opened
|
||||
InUse int
|
||||
Idle int
|
||||
WaitCount int64
|
||||
WaitDuration time.Duration
|
||||
MaxIdleClosed int64
|
||||
MaxLifetimeClosed int64
|
||||
}
|
||||
|
||||
// sqlConnection implements Connection for SQL databases (PostgreSQL, SQLite, MSSQL)
|
||||
@@ -415,6 +417,8 @@ func (c *sqlConnection) Stats() *ConnectionStats {
|
||||
if c.connected && c.provider != nil {
|
||||
if providerStats := c.provider.Stats(); providerStats != nil {
|
||||
stats.OpenConnections = providerStats.OpenConnections
|
||||
stats.MaxOpenConnections = providerStats.MaxOpenConnections
|
||||
stats.TotalOpened = providerStats.TotalOpened
|
||||
stats.InUse = providerStats.InUse
|
||||
stats.Idle = providerStats.Idle
|
||||
stats.WaitCount = providerStats.WaitCount
|
||||
|
||||
@@ -32,7 +32,7 @@ var (
|
||||
Name: "dbmanager_connection_pool_size",
|
||||
Help: "Current connection pool size",
|
||||
},
|
||||
[]string{"name", "type", "state"}, // state: open, idle, in_use
|
||||
[]string{"name", "type", "state"}, // state: open, idle, in_use, max
|
||||
)
|
||||
|
||||
// connectionWaitCount tracks how many times connections had to wait for availability
|
||||
@@ -71,6 +71,15 @@ var (
|
||||
[]string{"name", "type"},
|
||||
)
|
||||
|
||||
// connectionOpened tracks physical connections ever opened
|
||||
connectionOpened = promauto.NewCounterVec(
|
||||
prometheus.CounterOpts{
|
||||
Name: "dbmanager_connections_opened_total",
|
||||
Help: "Total physical connections opened by the pool",
|
||||
},
|
||||
[]string{"name", "type"},
|
||||
)
|
||||
|
||||
// connectionIdleClosed tracks connections closed due to max idle time
|
||||
connectionIdleClosed = promauto.NewCounterVec(
|
||||
prometheus.CounterOpts{
|
||||
@@ -115,6 +124,7 @@ func (m *connectionManager) PublishMetrics() {
|
||||
connectionPoolSize.WithLabelValues(name, string(connStats.Type), "open").Set(float64(connStats.OpenConnections))
|
||||
connectionPoolSize.WithLabelValues(name, string(connStats.Type), "idle").Set(float64(connStats.Idle))
|
||||
connectionPoolSize.WithLabelValues(name, string(connStats.Type), "in_use").Set(float64(connStats.InUse))
|
||||
connectionPoolSize.WithLabelValues(name, string(connStats.Type), "max").Set(float64(connStats.MaxOpenConnections))
|
||||
|
||||
// sql.DBStats values are cumulative, so add only the growth since
|
||||
// the last publish to keep these true counters.
|
||||
@@ -122,6 +132,7 @@ func (m *connectionManager) PublishMetrics() {
|
||||
connectionWaitCount.With(labels).Add(float64(connStats.WaitCount - prev.WaitCount))
|
||||
connectionWaitDuration.With(labels).Add((connStats.WaitDuration - prev.WaitDuration).Seconds())
|
||||
connectionLifetimeClosed.With(labels).Add(float64(connStats.MaxLifetimeClosed - prev.MaxLifetimeClosed))
|
||||
connectionOpened.With(labels).Add(float64(connStats.TotalOpened - prev.TotalOpened))
|
||||
connectionIdleClosed.With(labels).Add(float64(connStats.MaxIdleClosed - prev.MaxIdleClosed))
|
||||
}
|
||||
}
|
||||
@@ -152,7 +163,7 @@ func (p *publishedStats) swap(name string, cur *ConnectionStats) ConnectionStats
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
prev := p.last[name]
|
||||
if cur.WaitCount < prev.WaitCount || cur.MaxIdleClosed < prev.MaxIdleClosed || cur.MaxLifetimeClosed < prev.MaxLifetimeClosed {
|
||||
if cur.WaitCount < prev.WaitCount || cur.MaxIdleClosed < prev.MaxIdleClosed || cur.MaxLifetimeClosed < prev.MaxLifetimeClosed || cur.TotalOpened < prev.TotalOpened {
|
||||
prev = ConnectionStats{}
|
||||
}
|
||||
p.last[name] = *cur
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
package providers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"database/sql/driver"
|
||||
"sync/atomic"
|
||||
)
|
||||
|
||||
// countingConnector wraps a driver.Connector and counts every physical
|
||||
// connection it successfully dials. sql.DBStats has no such field, so this is
|
||||
// the only way to report the total number of connections ever opened.
|
||||
type countingConnector struct {
|
||||
driver.Connector
|
||||
opened *atomic.Int64
|
||||
}
|
||||
|
||||
func (c *countingConnector) Connect(ctx context.Context) (driver.Conn, error) {
|
||||
conn, err := c.Connector.Connect(ctx)
|
||||
if err == nil {
|
||||
c.opened.Add(1)
|
||||
}
|
||||
return conn, err
|
||||
}
|
||||
|
||||
// dsnConnector adapts a plain driver.Driver to driver.Connector.
|
||||
type dsnConnector struct {
|
||||
dsn string
|
||||
drv driver.Driver
|
||||
}
|
||||
|
||||
func (c dsnConnector) Connect(context.Context) (driver.Conn, error) { return c.drv.Open(c.dsn) }
|
||||
func (c dsnConnector) Driver() driver.Driver { return c.drv }
|
||||
|
||||
// openCounted is sql.Open with every dialled connection counted in opened.
|
||||
func openCounted(driverName, dsn string, opened *atomic.Int64) (*sql.DB, error) {
|
||||
probe, err := sql.Open(driverName, dsn)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
drv := probe.Driver()
|
||||
probe.Close() //nolint:gosec // G104: probe handle never dialled
|
||||
|
||||
var connector driver.Connector = dsnConnector{dsn: dsn, drv: drv}
|
||||
if dc, ok := drv.(driver.DriverContext); ok {
|
||||
if connector, err = dc.OpenConnector(dsn); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return sql.OpenDB(&countingConnector{Connector: connector, opened: opened}), nil
|
||||
}
|
||||
@@ -112,6 +112,12 @@ func (p *ExistingDBProvider) Stats() *ConnectionStats {
|
||||
if p.db != nil {
|
||||
dbStats := p.db.Stats()
|
||||
stats.OpenConnections = dbStats.OpenConnections
|
||||
stats.MaxOpenConnections = dbStats.MaxOpenConnections
|
||||
// The pool was opened outside dbmanager so dials cannot be counted.
|
||||
// Open plus every connection database/sql retired for idle or
|
||||
// lifetime limits is a close lower bound (it misses connections
|
||||
// dropped as broken).
|
||||
stats.TotalOpened = int64(dbStats.OpenConnections) + dbStats.MaxIdleClosed + dbStats.MaxIdleTimeClosed + dbStats.MaxLifetimeClosed
|
||||
stats.InUse = dbStats.InUse
|
||||
stats.Idle = dbStats.Idle
|
||||
stats.WaitCount = dbStats.WaitCount
|
||||
|
||||
@@ -3,9 +3,11 @@ package providers
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
_ "github.com/glebarez/sqlite"
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
)
|
||||
|
||||
@@ -159,6 +161,10 @@ func TestExistingDBProvider_Stats(t *testing.T) {
|
||||
t.Errorf("Expected stats.Type to be 'sql', got '%s'", stats.Type)
|
||||
}
|
||||
|
||||
if stats.MaxOpenConnections != 10 {
|
||||
t.Errorf("Expected stats.MaxOpenConnections to be 10, got %d", stats.MaxOpenConnections)
|
||||
}
|
||||
|
||||
if !stats.Connected {
|
||||
t.Error("Expected stats.Connected to be true")
|
||||
}
|
||||
@@ -192,3 +198,35 @@ func TestExistingDBProvider_Close_NilDB(t *testing.T) {
|
||||
t.Errorf("Expected Close to succeed with nil database, got error: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenCountedCountsDials(t *testing.T) {
|
||||
var opened atomic.Int64
|
||||
db, err := openCounted("sqlite", ":memory:", &opened)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
db.SetMaxOpenConns(1)
|
||||
|
||||
if err := db.PingContext(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := opened.Load(); got != 1 {
|
||||
t.Fatalf("opened = %d after first ping, want 1", got)
|
||||
}
|
||||
// Reusing the pooled connection must not count as a new dial.
|
||||
if err := db.PingContext(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := opened.Load(); got != 1 {
|
||||
t.Fatalf("opened = %d after reuse, want 1", got)
|
||||
}
|
||||
// Dropping idle connections forces a fresh dial.
|
||||
db.SetMaxIdleConns(0)
|
||||
if err := db.PingContext(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := opened.Load(); got != 2 {
|
||||
t.Fatalf("opened = %d after redial, want 2", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
_ "github.com/microsoft/go-mssqldb" // MSSQL driver
|
||||
@@ -16,6 +17,7 @@ import (
|
||||
type MSSQLProvider struct {
|
||||
db *sql.DB
|
||||
config ConnectionConfig
|
||||
opened atomic.Int64
|
||||
}
|
||||
|
||||
// NewMSSQLProvider creates a new MSSQL provider
|
||||
@@ -52,7 +54,7 @@ func (p *MSSQLProvider) Connect(ctx context.Context, cfg ConnectionConfig) error
|
||||
}
|
||||
|
||||
// Open database connection
|
||||
db, err = sql.Open("sqlserver", dsn)
|
||||
db, err = openCounted("sqlserver", dsn, &p.opened)
|
||||
if err != nil {
|
||||
lastErr = err
|
||||
if cfg.GetEnableLogging() {
|
||||
@@ -169,15 +171,17 @@ func (p *MSSQLProvider) Stats() *ConnectionStats {
|
||||
stats := p.db.Stats()
|
||||
|
||||
return &ConnectionStats{
|
||||
Name: p.config.GetName(),
|
||||
Type: "mssql",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
Name: p.config.GetName(),
|
||||
Type: "mssql",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
MaxOpenConnections: stats.MaxOpenConnections,
|
||||
TotalOpened: p.opened.Load(),
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"go.mongodb.org/mongo-driver/mongo"
|
||||
@@ -21,6 +22,7 @@ type PostgresProvider struct {
|
||||
config ConnectionConfig
|
||||
listener *PostgresListener
|
||||
mu sync.Mutex
|
||||
opened atomic.Int64
|
||||
}
|
||||
|
||||
// NewPostgresProvider creates a new PostgreSQL provider
|
||||
@@ -38,7 +40,7 @@ func (p *PostgresProvider) Connect(ctx context.Context, cfg ConnectionConfig) er
|
||||
// The connector and *sql.DB are created once; the pool is never closed to
|
||||
// recover from errors (see Refresh).
|
||||
connector := newPGConnector(connCfg)
|
||||
db := sql.OpenDB(connector)
|
||||
db := sql.OpenDB(&countingConnector{Connector: connector, opened: &p.opened})
|
||||
|
||||
// Connect with retry logic
|
||||
var lastErr error
|
||||
@@ -201,16 +203,18 @@ func (p *PostgresProvider) Stats() *ConnectionStats {
|
||||
stats := p.db.Stats()
|
||||
|
||||
return &ConnectionStats{
|
||||
Name: p.config.GetName(),
|
||||
Type: "postgres",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
Name: p.config.GetName(),
|
||||
Type: "postgres",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
MaxOpenConnections: stats.MaxOpenConnections,
|
||||
TotalOpened: p.opened.Load(),
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -27,13 +27,15 @@ type ConnectionStats struct {
|
||||
HealthCheckStatus string
|
||||
|
||||
// SQL connection pool stats
|
||||
OpenConnections int
|
||||
InUse int
|
||||
Idle int
|
||||
WaitCount int64
|
||||
WaitDuration time.Duration
|
||||
MaxIdleClosed int64
|
||||
MaxLifetimeClosed int64
|
||||
OpenConnections int
|
||||
MaxOpenConnections int
|
||||
TotalOpened int64 // physical connections ever dialled (0 when unknown)
|
||||
InUse int
|
||||
Idle int
|
||||
WaitCount int64
|
||||
WaitDuration time.Duration
|
||||
MaxIdleClosed int64
|
||||
MaxLifetimeClosed int64
|
||||
}
|
||||
|
||||
// ConnectionConfig is a minimal interface for configuration
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
_ "github.com/glebarez/sqlite" // Pure Go SQLite driver
|
||||
@@ -19,6 +20,7 @@ type SQLiteProvider struct {
|
||||
db *sql.DB
|
||||
dbMu sync.RWMutex
|
||||
config ConnectionConfig
|
||||
opened atomic.Int64
|
||||
}
|
||||
|
||||
// NewSQLiteProvider creates a new SQLite provider
|
||||
@@ -51,7 +53,7 @@ func (p *SQLiteProvider) Connect(ctx context.Context, cfg ConnectionConfig) erro
|
||||
}
|
||||
|
||||
// Open database connection
|
||||
db, err := sql.Open("sqlite", dsn)
|
||||
db, err := openCounted("sqlite", dsn, &p.opened)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to open SQLite connection: %w", err)
|
||||
}
|
||||
@@ -178,15 +180,17 @@ func (p *SQLiteProvider) Stats() *ConnectionStats {
|
||||
stats := p.db.Stats()
|
||||
|
||||
return &ConnectionStats{
|
||||
Name: p.config.GetName(),
|
||||
Type: "sqlite",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
Name: p.config.GetName(),
|
||||
Type: "sqlite",
|
||||
Connected: true,
|
||||
OpenConnections: stats.OpenConnections,
|
||||
MaxOpenConnections: stats.MaxOpenConnections,
|
||||
TotalOpened: p.opened.Load(),
|
||||
InUse: stats.InUse,
|
||||
Idle: stats.Idle,
|
||||
WaitCount: stats.WaitCount,
|
||||
WaitDuration: stats.WaitDuration,
|
||||
MaxIdleClosed: stats.MaxIdleClosed,
|
||||
MaxLifetimeClosed: stats.MaxLifetimeClosed,
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user