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
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.
351 lines
10 KiB
Go
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
|
|
}
|