/
rustwizard
/
pgtrace
Обзор
Документация
Войти
/
rustwizard
/
pgtrace
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
internal/pgstat/pgstat.go
303 строки
7 KB
Rust Wizard
feat: backend_type labels in timeline; overhead benchmark script
07 авг 2026, 14:48
Верифицирован
07 авг 2026, 14:48
0932132
Код
Авторство
О чём код?
// Package pgstat maintains a pid -> query cache backed by pg_stat_activity. package pgstat import ( "context" "fmt" "log/slog" "maps" "strconv" "sync" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) const ( activityQuery = `SELECT pid, query, state, backend_type FROM pg_stat_activity` activityQueryPG14 = `SELECT pid, query_id, query, state, backend_type FROM pg_stat_activity` stmtsQuery = `SELECT queryid, query FROM pg_stat_statements` extensionQuery = `SELECT EXISTS (SELECT 1 FROM pg_extension WHERE extname = 'pg_stat_statements')` serverVersionQuery = `SELECT current_setting('server_version_num')::int` ) // BackendType values from pg_stat_activity.backend_type. const ( BackendClient = "client backend" BackendAutovacuum = "autovacuum worker" BackendCheckpointer = "checkpointer" BackendBgWriter = "bgwriter" BackendWalWriter = "walwriter" ) // Activity describes what a backend pid is currently doing. type Activity struct { PID int32 Query string QueryID int64 // queryid from pg_stat_activity (0 when unavailable) State string BackendType string // pg_stat_activity.backend_type ("" when unavailable) } // IsBackground reports whether the backend is a server background process // rather than a client session. func (a Activity) IsBackground() bool { switch a.BackendType { case "", BackendClient: return false default: return true } } // Cache maps postgres backend pids to their current query, and queryids to // normalized query text from pg_stat_statements when the extension is present. type Cache struct { pool *pgxpool.Pool pg14 bool // server >= 14: pg_stat_activity has a query_id column stmts bool // pg_stat_statements extension is installed and readable mu sync.RWMutex m map[int32]Activity texts map[int64]string // queryid -> normalized query } // NewCache connects to Postgres, probes for pg_stat_statements support and // returns an empty cache. func NewCache(ctx context.Context, dsn string) (*Cache, error) { pool, err := pgxpool.New(ctx, dsn) if err != nil { return nil, fmt.Errorf("connect: %w", err) } if err := pool.Ping(ctx); err != nil { pool.Close() return nil, fmt.Errorf("ping: %w", err) } var version int if err := pool.QueryRow(ctx, serverVersionQuery).Scan(&version); err != nil { pool.Close() return nil, fmt.Errorf("detect server version: %w", err) } c := &Cache{ pool: pool, pg14: version >= 140000, m: make(map[int32]Activity), } var hasStmts bool if err := pool.QueryRow(ctx, extensionQuery).Scan(&hasStmts); err != nil { slog.Warn("pg_stat_statements detection failed", "err", err) } else { c.stmts = hasStmts if hasStmts { c.texts = make(map[int64]string) } } slog.Info("pgstat cache", "server_pg14", c.pg14, "pg_stat_statements", c.stmts) return c, nil } // Lookup returns the cached activity for pid, refreshing the whole cache on // a miss. Returns a zero Activity if the pid is unknown. func (c *Cache) Lookup(ctx context.Context, pid int32) Activity { c.mu.RLock() a, ok := c.m[pid] c.mu.RUnlock() if ok { return a } if err := c.Refresh(ctx); err != nil { slog.Warn("pg_stat_activity refresh", "err", err) return Activity{} } c.mu.RLock() a = c.m[pid] c.mu.RUnlock() return a } // Refresh reloads the whole pg_stat_activity snapshot and, when // pg_stat_statements is available, the queryid -> normalized query map. func (c *Cache) Refresh(ctx context.Context) error { q := activityQuery if c.pg14 { q = activityQueryPG14 } rows, err := c.pool.Query(ctx, q) if err != nil { return fmt.Errorf("query pg_stat_activity: %w", err) } fresh := make(map[int32]Activity) for rows.Next() { a, err := scanActivity(rows, c.pg14) if err != nil { rows.Close() return fmt.Errorf("scan pg_stat_activity: %w", err) } fresh[a.PID] = a } rows.Close() if err := rows.Err(); err != nil { return fmt.Errorf("iterate pg_stat_activity: %w", err) } if c.stmts { if err := c.refreshTexts(ctx); err != nil { slog.Warn("pg_stat_statements refresh", "err", err) } } c.mu.Lock() c.m = fresh c.mu.Unlock() return nil } // scanActivity scans one pg_stat_activity row into an Activity. func scanActivity(rows pgx.Rows, pg14 bool) (Activity, error) { var ( pid int32 queryID *int64 query, state, backendType *string ) if pg14 { if err := rows.Scan(&pid, &queryID, &query, &state, &backendType); err != nil { return Activity{}, err } } else { if err := rows.Scan(&pid, &query, &state, &backendType); err != nil { return Activity{}, err } } a := Activity{PID: pid} if queryID != nil { a.QueryID = *queryID } if query != nil { a.Query = *query } if state != nil { a.State = *state } if backendType != nil { a.BackendType = *backendType } return a, nil } // refreshTexts reloads the queryid -> normalized query map from // pg_stat_statements. On failure the previous map is kept. func (c *Cache) refreshTexts(ctx context.Context) error { rows, err := c.pool.Query(ctx, stmtsQuery) if err != nil { return fmt.Errorf("query pg_stat_statements: %w", err) } defer rows.Close() texts := make(map[int64]string) for rows.Next() { var ( queryID int64 query string ) if err := rows.Scan(&queryID, &query); err != nil { return fmt.Errorf("scan pg_stat_statements: %w", err) } texts[queryID] = query } if err := rows.Err(); err != nil { return fmt.Errorf("iterate pg_stat_statements: %w", err) } c.mu.Lock() c.texts = texts c.mu.Unlock() return nil } // QueryText returns the normalized query text for queryID, or "" when // pg_stat_statements is unavailable or the queryid is not tracked yet. func (c *Cache) QueryText(queryID int64) string { c.mu.RLock() t := c.texts[queryID] c.mu.RUnlock() return t } // QueryTexts returns a copy of the queryid -> normalized query mapping. func (c *Cache) QueryTexts() map[int64]string { c.mu.RLock() defer c.mu.RUnlock() out := make(map[int64]string, len(c.texts)) maps.Copy(out, c.texts) return out } // Resolve returns the stable aggregation key for act: the normalized // template text when a queryid is known and text is available, otherwise // the raw query text, or "" when no query is known. func (c *Cache) Resolve(act Activity) string { if act.QueryID != 0 { if t := c.QueryText(act.QueryID); t != "" { return t } return "queryid=" + strconv.FormatInt(act.QueryID, 10) } return act.Query } // StartRefresher periodically refreshes the cache until ctx is cancelled. func (c *Cache) StartRefresher(ctx context.Context, interval time.Duration) { go func() { ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if err := c.Refresh(ctx); err != nil { slog.Warn("pg_stat_activity refresh", "err", err) } } } }() } // Close releases the connection pool. func (c *Cache) Close() { c.pool.Close() }