Files
Denozordec 9efa3bbc8a
CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Has been skipped
CI / web (push) Successful in 30s
CI / go (push) Failing after 24s
CI / bird2 (push) Has been skipped
CI / release (push) Has been skipped
feat(db): enhance PostgreSQL statistics monitoring and error handling
Updated the PostgreSQL monitoring service to improve handling of `pg_stat_statements` availability. Introduced a new method to check if the extension is queryable and updated the response structure to include availability status and hints. Enhanced the documentation to clarify the requirements for enabling `pg_stat_statements`. Adjusted related components to reflect these changes, ensuring better user feedback in the monitoring interface.
2026-06-01 14:15:38 +07:00

351 lines
10 KiB
Go

package pgmonitor
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
func clampLimit(limit, def, max int) int {
if limit <= 0 {
return def
}
if limit > max {
return max
}
return limit
}
func (s *Service) fetchOverview(ctx context.Context) (Overview, error) {
now := time.Now().UTC()
out := Overview{CollectedAt: now}
var active, idle, total, maxConn int
err := s.pool.QueryRow(ctx, `
SELECT
count(*) FILTER (WHERE state = 'active'),
count(*) FILTER (WHERE state = 'idle'),
count(*),
(SELECT setting::int FROM pg_settings WHERE name = 'max_connections')
FROM pg_stat_activity
WHERE datname = current_database()`).Scan(&active, &idle, &total, &maxConn)
if err != nil {
return out, fmt.Errorf("pgmonitor: connections: %w", err)
}
out.Connections = Connections{Active: active, Idle: idle, Total: total, MaxConnections: maxConn}
var cachePct *float64
err = s.pool.QueryRow(ctx, `
SELECT numbackends, xact_commit, xact_rollback, deadlocks, blks_hit, blks_read,
CASE WHEN blks_hit + blks_read > 0
THEN round(100.0 * blks_hit::numeric / (blks_hit + blks_read), 2) END
FROM pg_stat_database WHERE datname = current_database()`).Scan(
&out.Database.Backends,
&out.Database.XactCommit,
&out.Database.XactRollback,
&out.Database.Deadlocks,
&out.Database.BlksHit,
&out.Database.BlksRead,
&cachePct,
)
if err != nil {
return out, fmt.Errorf("pgmonitor: database stats: %w", err)
}
if cachePct != nil {
out.Database.CacheHitPct = *cachePct
}
_ = s.pool.QueryRow(ctx, `
SELECT checkpoints_timed, checkpoints_req, buffers_checkpoint, buffers_clean,
maxwritten_clean, buffers_backend, buffers_alloc
FROM pg_stat_bgwriter`).Scan(
&out.Bgwriter.CheckpointsTimed,
&out.Bgwriter.CheckpointsReq,
&out.Bgwriter.BuffersCheckpoint,
&out.Bgwriter.BuffersClean,
&out.Bgwriter.MaxWrittenClean,
&out.Bgwriter.BuffersBackend,
&out.Bgwriter.BuffersAlloc,
)
_ = s.pool.QueryRow(ctx, `SELECT pg_database_size(current_database())`).Scan(&out.SizeBytes)
_ = s.pool.QueryRow(ctx, `
SELECT
(SELECT setting FROM pg_settings WHERE name = 'shared_buffers'),
(SELECT setting FROM pg_settings WHERE name = 'work_mem'),
(SELECT setting FROM pg_settings WHERE name = 'effective_cache_size')`).Scan(
&out.MemorySettings.SharedBuffers,
&out.MemorySettings.WorkMem,
&out.MemorySettings.EffectiveCacheSize,
)
rows, err := s.pool.Query(ctx, `
SELECT client_addr::text, state, sync_state,
EXTRACT(EPOCH FROM COALESCE(write_lag, flush_lag, replay_lag)) * 1000
FROM pg_stat_replication`)
if err == nil {
defer rows.Close()
for rows.Next() {
var peer ReplicationPeer
var lagMs *float64
if err := rows.Scan(&peer.ClientAddr, &peer.State, &peer.SyncState, &lagMs); err != nil {
continue
}
if lagMs != nil {
v := int64(*lagMs)
peer.LagMs = &v
}
out.Replication = append(out.Replication, peer)
}
}
out.StatementsEnabled = s.statementsQueryable(ctx)
return out, nil
}
func queryLocks(ctx context.Context, pool *pgxpool.Pool) ([]LockRow, error) {
rows, err := pool.Query(ctx, `
SELECT l.locktype, l.mode, l.granted, a.pid, COALESCE(a.usename, ''),
COALESCE(a.state, ''), COALESCE(left(a.query, 300), ''),
NOT l.granted AS blocked
FROM pg_locks l
JOIN pg_stat_activity a ON a.pid = l.pid
WHERE a.datname = current_database()
AND (NOT l.granted OR l.mode LIKE '%Exclusive%')
ORDER BY l.granted ASC, a.query_start NULLS LAST
LIMIT 200`)
if err != nil {
return nil, fmt.Errorf("pgmonitor: locks: %w", err)
}
defer rows.Close()
var out []LockRow
for rows.Next() {
var r LockRow
if err := rows.Scan(&r.Locktype, &r.Mode, &r.Granted, &r.PID, &r.User, &r.State, &r.Query, &r.Blocked); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
func queryTables(ctx context.Context, pool *pgxpool.Pool, limit int) ([]TableStat, error) {
limit = clampLimit(limit, 20, 100)
rows, err := pool.Query(ctx, `
SELECT t.relname,
pg_total_relation_size(t.relid),
s.heap_blks_read, s.heap_blks_hit,
t.idx_scan, t.seq_scan, t.n_dead_tup, t.last_autovacuum,
CASE WHEN t.n_live_tup + t.n_dead_tup > 0
THEN round(t.n_dead_tup::numeric / (t.n_live_tup + t.n_dead_tup), 4)
ELSE 0 END
FROM pg_statio_user_tables s
JOIN pg_stat_user_tables t ON t.relid = s.relid
WHERE t.schemaname = 'public'
ORDER BY pg_total_relation_size(t.relid) DESC
LIMIT $1`, limit)
if err != nil {
return nil, fmt.Errorf("pgmonitor: tables: %w", err)
}
defer rows.Close()
var out []TableStat
for rows.Next() {
var r TableStat
var last *time.Time
if err := rows.Scan(&r.Relname, &r.TotalBytes, &r.HeapBlksRead, &r.HeapBlksHit,
&r.IdxScan, &r.SeqScan, &r.DeadTuples, &last, &r.BloatRatio); err != nil {
return nil, err
}
r.LastAutovacuum = last
out = append(out, r)
}
return out, rows.Err()
}
// TopQueries loads from pg_stat_statements when available.
func (s *Service) TopQueries(ctx context.Context, limit int) (QueriesResponse, error) {
if s == nil || s.pool == nil {
return QueriesResponse{}, errors.New("pgmonitor: postgres not configured")
}
limit = clampLimit(limit, 20, 100)
now := time.Now().UTC()
if snap, ok, err := s.loadSnapshot(ctx, "slow_queries", 15*time.Minute); err == nil && ok {
var items []QueryStat
if err := decodePayload(snap.Payload, &items); err == nil {
return QueriesResponse{
CollectedAt: snap.CollectedAt,
Source: "snapshot",
Items: items,
StatementsAvailable: true,
}, nil
}
}
if !s.statementsQueryable(ctx) {
return queriesUnavailable(now), nil
}
items, err := queryTopStatements(ctx, s.pool, limit)
if err != nil {
if isPgStatStatementsUnavailable(err) {
s.markStatementsUnavailable()
return queriesUnavailable(now), nil
}
return QueriesResponse{}, err
}
return QueriesResponse{
CollectedAt: now,
Source: "live",
Items: items,
StatementsAvailable: true,
}, nil
}
func queriesUnavailable(at time.Time) QueriesResponse {
return QueriesResponse{
CollectedAt: at,
Source: "unavailable",
Items: nil,
StatementsAvailable: false,
StatementsHint: statementsUnavailableHint,
}
}
const statementsUnavailableHint = "pg_stat_statements requires shared_preload_libraries and PostgreSQL restart (see docs/db-diagnostics.md)"
// statementsQueryable returns true only when pg_stat_statements can be queried (not merely installed).
func (s *Service) statementsQueryable(ctx context.Context) bool {
if s == nil || s.pool == nil {
return false
}
if v, ok := s.cache.get("stmt_queryable"); ok {
if b, ok := v.(bool); ok {
return b
}
}
ok := probePgStatStatements(ctx, s.pool)
s.cache.set("stmt_queryable", ok)
return ok
}
func (s *Service) markStatementsUnavailable() {
s.cache.set("stmt_queryable", false)
}
func probePgStatStatements(ctx context.Context, pool *pgxpool.Pool) bool {
var dummy int64
err := pool.QueryRow(ctx, `
SELECT COALESCE(SUM(calls), 0)::bigint FROM pg_stat_statements LIMIT 1`).Scan(&dummy)
if err == nil {
return true
}
return !isPgStatStatementsUnavailable(err)
}
func queryTopStatements(ctx context.Context, pool *pgxpool.Pool, limit int) ([]QueryStat, error) {
rows, err := pool.Query(ctx, `
SELECT queryid, left(query, 500), calls, total_exec_time, mean_exec_time, rows
FROM pg_stat_statements
WHERE dbid = (SELECT oid FROM pg_database WHERE datname = current_database())
ORDER BY mean_exec_time DESC
LIMIT $1`, limit)
if err != nil {
if isPgStatStatementsUnavailable(err) {
return nil, nil
}
return nil, fmt.Errorf("pgmonitor: pg_stat_statements: %w", err)
}
defer rows.Close()
var out []QueryStat
for rows.Next() {
var r QueryStat
if err := rows.Scan(&r.QueryID, &r.Query, &r.Calls, &r.TotalExecMs, &r.MeanExecMs, &r.Rows); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
// isPgStatStatementsUnavailable reports extension missing or not loaded via shared_preload_libraries.
func isPgStatStatementsUnavailable(err error) bool {
if err == nil {
return false
}
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) {
switch pgErr.Code {
case "42P01", "42704", "55000":
return true
}
msg := strings.ToLower(pgErr.Message)
if strings.Contains(msg, "shared_preload_libraries") || strings.Contains(msg, "pg_stat_statements") {
return true
}
}
low := strings.ToLower(err.Error())
return strings.Contains(low, "shared_preload_libraries") || strings.Contains(low, "pg_stat_statements")
}
func isSafeIdent(name string) bool {
if name == "" {
return true
}
for _, r := range name {
if (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') || r == '_' {
continue
}
return false
}
return true
}
// ExecMaintenance runs VACUUM/ANALYZE/REINDEX with optional dry-run (returns SQL executed or planned).
func ExecMaintenance(ctx context.Context, pool *pgxpool.Pool, kind, table string, dryRun bool) (detail map[string]any, err error) {
if pool == nil {
return nil, errors.New("pgmonitor: postgres not configured")
}
table = strings.TrimSpace(table)
if table != "" && !isSafeIdent(table) {
return nil, errors.New("pgmonitor: invalid table name")
}
qual := ""
if table != "" {
qual = " " + pgx.Identifier{table}.Sanitize()
}
var sql string
switch kind {
case "vacuum":
sql = "VACUUM" + qual
case "vacuum_analyze":
sql = "VACUUM ANALYZE" + qual
case "analyze":
sql = "ANALYZE" + qual
case "reindex":
if table == "" {
return nil, errors.New("pgmonitor: reindex requires table")
}
sql = "REINDEX TABLE" + qual
default:
return nil, fmt.Errorf("pgmonitor: unknown maintenance kind %q", kind)
}
detail = map[string]any{"sql": sql, "dry_run": dryRun}
if dryRun {
return detail, nil
}
_, err = pool.Exec(ctx, sql)
if err != nil {
return detail, fmt.Errorf("pgmonitor: %s: %w", kind, err)
}
detail["executed"] = true
return detail, nil
}