Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ The following collectors are configurable:
| Name | Description | Enabled by default |
|------------------|-----------------------------------------------------------------------|--------------------|
| `explain_plans` | Collect query explain plans. | yes |
| `table_stats` | Collect table-level scan statistics. | no |
| `index_stats` | Collect per-index usage statistics. | no |
| `query_details` | Collect queries information. | yes |
| `query_samples` | Collect query samples and wait events information. | yes |
| `schema_details` | Collect schemas, tables, and columns from PostgreSQL system catalogs. | yes |
Expand Down
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package collector

import (
"context"
"database/sql"
"errors"
"fmt"
"regexp"
)

Expand All @@ -13,6 +15,14 @@ var defaultDbConnectionFactory = func(dsn string) (*sql.DB, error) {
return sql.Open("postgres", dsn)
}

// selectAllDatabases makes use of the initial DB connection to discover other databases on the same Postgres instance
const selectAllDatabases = `
SELECT datname
FROM pg_database
WHERE datistemplate = false
AND has_database_privilege(datname, 'CONNECT')
AND datname NOT IN %s`

// replaceDatabaseNameInDSN safely replaces the database name in a PostgreSQL DSN
// using regex to ensure only the database name portion is replaced, not other occurrences
func replaceDatabaseNameInDSN(dsn, newDatabaseName string) (string, error) {
Expand All @@ -30,3 +40,74 @@ func replaceDatabaseNameInDSN(dsn, newDatabaseName string) (string, error) {
newDSN := matches[1] + newDatabaseName + matches[3]
return newDSN, nil
}

// databaseNameFromDSN extracts the database name a DSN already points to, so
// callers fanning out per-database can tell whether a target database is the
// one an existing connection already uses.
func databaseNameFromDSN(dsn string) (string, error) {
matches := dsnParseRegex.FindStringSubmatch(dsn)
if len(matches) < 4 {
return "", errors.New("failed to parse DSN for database name")
}
return matches[2], nil
}

// discoverDatabases lists databases the current connection can reach, via
// pg_database (readable from any single connection) -- used to fan out
// per-database connections, since most stat views only report on the
// database a connection is actually established to.
func discoverDatabases(ctx context.Context, conn *sql.DB, excludeDatabases []string) ([]string, error) {
query := fmt.Sprintf(selectAllDatabases, buildExcludedDatabasesClause(excludeDatabases))
rows, err := conn.QueryContext(ctx, query)
if err != nil {
return nil, fmt.Errorf("failed to discover databases: %w", err)
}
defer rows.Close()

var databases []string
for rows.Next() {
var datname string
if err := rows.Scan(&datname); err != nil {
return nil, fmt.Errorf("failed to scan database name: %w", err)
}
databases = append(databases, datname)
}

if err := rows.Err(); err != nil {
return nil, fmt.Errorf("error iterating database rows: %w", err)
}

return databases, nil
}

// connectToDatabase opens a connection to dbName by rewriting dsn. If dbName
// is already what dsn (and so initial) points to, it reuses initial instead
// of opening a redundant connection -- sql.Open never returns something
// pointer-equal to an existing *sql.DB, so this has to be checked by name up
// front, not via "conn != initial" after the fact. closeFn closes the
// connection unless it's initial.
func connectToDatabase(dsn, dbName string, factory databaseConnectionFactory, initial *sql.DB) (conn *sql.DB, closeFn func(), err error) {
noopClose := func() {}

if currentDBName, err := databaseNameFromDSN(dsn); err == nil && currentDBName == dbName {
return initial, noopClose, nil
}

databaseDSN, err := replaceDatabaseNameInDSN(dsn, dbName)
if err != nil {
return nil, nil, fmt.Errorf("failed to create DSN for database %s: %w", dbName, err)
}

conn, err = factory(databaseDSN)
if err != nil {
return nil, nil, fmt.Errorf("failed to create connection to database %s: %w", dbName, err)
}

closeFn = func() {
if conn != initial {
conn.Close()
}
}

return conn, closeFn, nil
}
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package collector

import (
"database/sql"
"testing"

sqlmock "github.com/DATA-DOG/go-sqlmock"
"github.com/stretchr/testify/require"
)

Expand Down Expand Up @@ -96,3 +98,47 @@ func TestReplaceDatabaseNameInDSN(t *testing.T) {
})
}
}

// TestConnectToDatabaseReusesInitialConnection guards against connectToDatabase
// opening a redundant connection for the database initial already points to;
// see the comment on connectToDatabase for why a bare "conn != initial" check
// can't catch this.
func TestConnectToDatabaseReusesInitialConnection(t *testing.T) {
initial, _, err := sqlmock.New()
require.NoError(t, err)
defer initial.Close()

newDB, _, err := sqlmock.New()
require.NoError(t, err)
defer newDB.Close()

t.Run("same database as the DSN: reuses initial, never calls factory", func(t *testing.T) {
factoryCalls := 0
factory := func(dsn string) (*sql.DB, error) {
factoryCalls++
return newDB, nil
}

conn, closeFn, err := connectToDatabase("postgres://user:pass@localhost:5432/books_store", "books_store", factory, initial)
require.NoError(t, err)
require.Same(t, initial, conn)
require.Equal(t, 0, factoryCalls)
closeFn() // must not close initial

require.NoError(t, initial.PingContext(t.Context())) // still usable
})

t.Run("different database: opens a new connection via factory", func(t *testing.T) {
factoryCalls := 0
factory := func(dsn string) (*sql.DB, error) {
factoryCalls++
return newDB, nil
}

conn, closeFn, err := connectToDatabase("postgres://user:pass@localhost:5432/postgres", "books_store", factory, initial)
require.NoError(t, err)
require.Same(t, newDB, conn)
require.Equal(t, 1, factoryCalls)
closeFn()
})
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
package collector

import (
"context"
"database/sql"
"log/slog"
"strconv"

"github.com/prometheus/client_golang/prometheus"
"go.uber.org/atomic"
)

// IndexStatsCollector emits per-index usage counters from pg_stat_user_indexes,
// scoped to every database the connection can reach rather than only the one
// named in the DSN.
const IndexStatsCollector = "index_stats"

const selectIndexUsageStats = `
SELECT
s.schemaname,
s.relname,
s.indexrelname,
s.idx_scan,
i.indisprimary,
i.indisunique,
i.indpred IS NOT NULL AS is_partial,
pg_relation_size(s.indexrelid) AS index_size_bytes
FROM pg_stat_user_indexes s
JOIN pg_index i ON i.indexrelid = s.indexrelid`

var indexLabels = []string{labelDatname, "schemaname", "relname", "indexrelname"}
var indexSizeLabels = append(append([]string{}, indexLabels...), "is_primary", "is_unique", "is_partial")

var (
indexUsageIdxScanTotalDesc = prometheus.NewDesc(
prometheus.BuildFQName("database_observability", "pg_index_stats", "idx_scan_total"),
"Number of index scans initiated on this index",
indexLabels, nil,
)
indexSizeBytesDesc = prometheus.NewDesc(
prometheus.BuildFQName("database_observability", "pg_index_stats", "size_bytes"),
"Total disk space used by this index, in bytes, labeled with whether it backs the primary key or a unique constraint, or is partial",
indexSizeLabels, nil,
)
)

type IndexStatsArguments struct {
DB *sql.DB
DSN string
ExcludeDatabases []string
Registry *prometheus.Registry

Logger *slog.Logger

dbConnectionFactory databaseConnectionFactory
}

type IndexStats struct {
initialConnection *sql.DB
dbDSN string
dbConnectionFactory databaseConnectionFactory
excludeDatabases []string
registry *prometheus.Registry

logger *slog.Logger
running *atomic.Bool
}

func NewIndexStats(args IndexStatsArguments) (*IndexStats, error) {
factory := args.dbConnectionFactory
if factory == nil {
factory = defaultDbConnectionFactory
}

return &IndexStats{
initialConnection: args.DB,
dbDSN: args.DSN,
dbConnectionFactory: factory,
excludeDatabases: args.ExcludeDatabases,
registry: args.Registry,
logger: args.Logger.With("collector", IndexStatsCollector),
running: &atomic.Bool{},
}, nil
}

func (c *IndexStats) Name() string {
return IndexStatsCollector
}

func (c *IndexStats) Start(_ context.Context) error {
if err := c.registry.Register(c); err != nil {
return err
}
c.running.Store(true)
return nil
}

func (c *IndexStats) Stopped() bool {
return !c.running.Load()
}

func (c *IndexStats) Stop() {
c.registry.Unregister(c)
c.running.Store(false)
}

// Describe implements prometheus.Collector.
func (c *IndexStats) Describe(ch chan<- *prometheus.Desc) {
ch <- indexUsageIdxScanTotalDesc
ch <- indexSizeBytesDesc
}

// Collect implements prometheus.Collector. It runs synchronously at scrape
// time, fanning out to every database the connection can reach.
func (c *IndexStats) Collect(ch chan<- prometheus.Metric) {
ctx := context.Background()

databases, err := discoverDatabases(ctx, c.initialConnection, c.excludeDatabases)
if err != nil {
c.logger.Error("failed to discover databases", "err", err)
return
}

for _, dbName := range databases {
conn, closeConn, err := connectToDatabase(c.dbDSN, dbName, c.dbConnectionFactory, c.initialConnection)
if err != nil {
c.logger.Error("failed to connect to database", "datname", dbName, "err", err)
continue
}

c.collectIndexUsageStats(ctx, dbName, conn, ch)

closeConn()
}
}

func (c *IndexStats) collectIndexUsageStats(ctx context.Context, dbName string, conn *sql.DB, ch chan<- prometheus.Metric) {
rows, err := conn.QueryContext(ctx, selectIndexUsageStats)
if err != nil {
c.logger.Error("failed to query pg_stat_user_indexes", "datname", dbName, "err", err)
return
}
defer rows.Close()

for rows.Next() {
var schemaname, relname, indexrelname string
var idxScan, indexSizeBytes sql.NullInt64
var isPrimary, isUnique, isPartial bool

if err := rows.Scan(&schemaname, &relname, &indexrelname, &idxScan, &isPrimary, &isUnique, &isPartial, &indexSizeBytes); err != nil {
c.logger.Error("failed to scan pg_stat_user_indexes row", "datname", dbName, "err", err)
return
}

ch <- prometheus.MustNewConstMetric(indexUsageIdxScanTotalDesc, prometheus.CounterValue, float64(idxScan.Int64), dbName, schemaname, relname, indexrelname)
ch <- prometheus.MustNewConstMetric(indexSizeBytesDesc, prometheus.GaugeValue, float64(indexSizeBytes.Int64),
dbName, schemaname, relname, indexrelname,
strconv.FormatBool(isPrimary), strconv.FormatBool(isUnique), strconv.FormatBool(isPartial))
}

if err := rows.Err(); err != nil {
c.logger.Error("error iterating pg_stat_user_indexes rows", "datname", dbName, "err", err)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
package collector

import (
"database/sql"
"fmt"
"strings"
"testing"

sqlmock "github.com/DATA-DOG/go-sqlmock"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/testutil"
"github.com/stretchr/testify/require"

"github.com/grafana/alloy/internal/util"
)

func TestIndexStats(t *testing.T) {
db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual))
require.NoError(t, err)
defer db.Close()

registry := prometheus.NewRegistry()

c, err := NewIndexStats(IndexStatsArguments{
DB: db,
DSN: "postgres://user:pass@localhost:5432/books_store",
ExcludeDatabases: nil,
Registry: registry,
Logger: util.TestAlloyLogger(t).Slog(),
dbConnectionFactory: func(dsn string) (*sql.DB, error) {
return db, nil
},
})
require.NoError(t, err)

require.NoError(t, c.Start(t.Context()))
defer c.Stop()

mock.ExpectQuery(fmt.Sprintf(selectAllDatabases, exclusionClause)).WithoutArgs().RowsWillBeClosed().
WillReturnRows(sqlmock.NewRows([]string{"datname"}).AddRow("books_store"))

mock.ExpectQuery(selectIndexUsageStats).WithoutArgs().RowsWillBeClosed().
WillReturnRows(
sqlmock.NewRows([]string{"schemaname", "relname", "indexrelname", "idx_scan", "indisprimary", "indisunique", "is_partial", "index_size_bytes"}).
AddRow("public", "books", "books_pkey", 184000000, true, true, false, 65536).
AddRow("public", "books", "idx_books_title", 0, false, false, true, 32768),
)

expected := `
# HELP database_observability_pg_index_stats_idx_scan_total Number of index scans initiated on this index
# TYPE database_observability_pg_index_stats_idx_scan_total counter
database_observability_pg_index_stats_idx_scan_total{datname="books_store",indexrelname="books_pkey",relname="books",schemaname="public"} 1.84e+08
database_observability_pg_index_stats_idx_scan_total{datname="books_store",indexrelname="idx_books_title",relname="books",schemaname="public"} 0
# HELP database_observability_pg_index_stats_size_bytes Total disk space used by this index, in bytes, labeled with whether it backs the primary key or a unique constraint, or is partial
# TYPE database_observability_pg_index_stats_size_bytes gauge
database_observability_pg_index_stats_size_bytes{datname="books_store",indexrelname="books_pkey",is_partial="false",is_primary="true",is_unique="true",relname="books",schemaname="public"} 65536
database_observability_pg_index_stats_size_bytes{datname="books_store",indexrelname="idx_books_title",is_partial="true",is_primary="false",is_unique="false",relname="books",schemaname="public"} 32768
`

require.NoError(t, testutil.CollectAndCompare(registry, strings.NewReader(expected)))
require.NoError(t, mock.ExpectationsWereMet())
}
Loading
Loading