diff --git a/README.md b/README.md index 4d530e5..c876fc2 100644 --- a/README.md +++ b/README.md @@ -404,8 +404,12 @@ Global settings applied to all database instances: | `worker_idle_timeout_sec` | Worker idle timeout | `10` | | `notify_retry_seconds` | NOTIFY retry interval | `30s` | | `enable_debug` | Enable debug logging | `false` | -| `lease_seconds` | Job lease duration before reclaimable | - | -| `stale_job_recovery_sec` | Interval for reclaiming expired leases | - | +| `lease_seconds` | Job lease duration before reclaimable | `60` | +| `stale_job_recovery_sec` | Interval for reclaiming expired leases | `30` | +| `metrics_enabled` | Enable Prometheus and the embedded dashboard | `true` | +| `metrics_host` | Metrics HTTP bind host | `0.0.0.0` | +| `metrics_port` | Metrics HTTP bind port | `9469` | +| `queue_depth_poll_sec` | Queue depth collection interval | `15` | ## Development diff --git a/broker.example.yaml b/broker.example.yaml index c4a1215..678fc1a 100644 --- a/broker.example.yaml +++ b/broker.example.yaml @@ -41,6 +41,10 @@ broker: worker_idle_timeout_sec: 10 # Worker idle timeout notify_retry_seconds: 30s # LISTEN/NOTIFY retry interval enable_debug: false # Enable debug logging + metrics_enabled: true # Expose Prometheus and the embedded dashboard + metrics_host: 0.0.0.0 + metrics_port: 9469 + queue_depth_poll_sec: 15 # Logging settings logging: diff --git a/cmd/broker/db.go b/cmd/broker/db.go index f7e8a37..952fed7 100644 --- a/cmd/broker/db.go +++ b/cmd/broker/db.go @@ -355,7 +355,8 @@ func runDBList() error { return nil } fmt.Printf("%-20s %-20s %-6s %-16s %-16s %-10s\n", "NAME", "HOST", "PORT", "DATABASE", "USER", "STATUS") - for _, db := range dbs { + for i := range dbs { + db := &dbs[i] status := "enabled" if db.Disabled { status = "disabled" diff --git a/go.mod b/go.mod index d62e5d8..fb5bebf 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.26.0 require ( github.com/lib/pq v1.10.9 + github.com/prometheus/client_golang v1.20.5 github.com/spf13/cobra v1.10.2 github.com/spf13/viper v1.21.0 github.com/stretchr/testify v1.11.1 @@ -12,12 +13,19 @@ require ( ) require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/go-viper/mapstructure/v2 v2.4.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect + github.com/klauspost/compress v1.17.9 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/prometheus/client_model v0.6.1 // indirect + github.com/prometheus/common v0.55.0 // indirect + github.com/prometheus/procfs v0.15.1 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect github.com/spf13/afero v1.15.0 // indirect @@ -27,4 +35,5 @@ require ( go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/sys v0.48.0 // indirect golang.org/x/text v0.28.0 // indirect + google.golang.org/protobuf v1.34.2 // indirect ) diff --git a/go.sum b/go.sum index 9d1a74f..8c22bad 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,7 @@ +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -11,18 +15,32 @@ github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= +github.com/klauspost/compress v1.17.9 h1:6KIumPrER1LHsvBVuDa0r5xaG0Es51mhhB9BQB2qeMA= +github.com/klauspost/compress v1.17.9/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/lib/pq v1.10.9 h1:YXG7RB+JIjhP29X+OtkiDnYaXQwpS4JEWq7dtCCRUEw= github.com/lib/pq v1.10.9/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= -github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs= +github.com/prometheus/client_golang v1.20.5 h1:cxppBPuYhUnsO6yo/aoRol4L7q7UFfdm+bR9r+8l63Y= +github.com/prometheus/client_golang v1.20.5/go.mod h1:PIEt8X02hGcP8JWbeHyeZ53Y/jReSnHgO035n//V5WE= +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= +github.com/prometheus/common v0.55.0 h1:KEi6DK7lXW/m7Ig5i47x0vRzuBsHuvJdi5ee6Y3G1dc= +github.com/prometheus/common v0.55.0/go.mod h1:2SECS4xJG1kd8XF9IcM1gMX6510RAEL65zxzNImwdc8= +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= +github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ= +github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDcg+AAIFXc= github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik= @@ -51,8 +69,10 @@ golang.org/x/term v0.46.0 h1:3+OXuTbaKDgwk8jTi3aSLHRlmWqHEUDUtxnbFigO4YE= golang.org/x/term v0.46.0/go.mod h1:+K02xbkittuwc0Am4abfA3Fc+XRGXkvBXNO88NCXPoc= golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng= golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU= +google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg= +google.golang.org/protobuf v1.34.2/go.mod h1:qYOHts0dSfpeUzUFpOMr/WGzszTmLH+DiWniOlNbLDw= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15 h1:YR8cESwS4TdDjEe65xsg0ogRM/Nc3DYOhEAlW+xobZo= -gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/pkg/broker/broker.go b/pkg/broker/broker.go index 3501ea8..1027ba1 100644 --- a/pkg/broker/broker.go +++ b/pkg/broker/broker.go @@ -4,9 +4,11 @@ import ( "context" "fmt" "sync" + "time" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/config" + "git.warky.dev/wdevs/pgsql-broker/pkg/broker/metrics" ) // Broker manages multiple database instances @@ -17,6 +19,8 @@ type Broker struct { instances []*DatabaseInstance ctx context.Context cancel context.CancelFunc + metrics *metrics.Metrics + server *metrics.Server shutdown bool mu sync.RWMutex } @@ -33,6 +37,16 @@ func New(cfg *config.Config, logger adapter.Logger, version string) (*Broker, er ctx: ctx, cancel: cancel, } + if cfg.Broker.MetricsEnabled { + broker.metrics = metrics.New() + broker.metrics.SetDatabaseCount(len(cfg.Databases)) + addr := fmt.Sprintf("%s:%d", cfg.Broker.MetricsHost, cfg.Broker.MetricsPort) + server, err := metrics.NewServer(broker.metrics, addr, broker.logger) + if err != nil { + return nil, err + } + broker.server = server + } return broker, nil } @@ -54,7 +68,7 @@ func (b *Broker) Start() error { dbAdapter := adapter.NewPostgresAdapter(dbCfg.ToPostgresConfig(), b.logger) // Create database instance - instance, err := NewDatabaseInstance(b.config, dbCfg, dbAdapter, b.logger, b.version, b.ctx) + instance, err := NewDatabaseInstance(b.config, dbCfg, dbAdapter, b.logger, b.version, b.ctx, b.metrics) if err != nil { // Stop any already-started instances b.stopInstances() @@ -77,6 +91,9 @@ func (b *Broker) Start() error { } b.logger.Info("broker started successfully", "database_instances", len(b.instances)) + if b.server != nil { + b.server.Start() + } return nil } @@ -97,6 +114,13 @@ func (b *Broker) Stop() error { // Stop all instances b.stopInstances() + if b.server != nil { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := b.server.Stop(ctx); err != nil { + b.logger.Error("failed to stop metrics server", "error", err) + } + } b.logger.Info("broker stopped") return nil diff --git a/pkg/broker/config/config.go b/pkg/broker/config/config.go index dfea8aa..48d9df8 100644 --- a/pkg/broker/config/config.go +++ b/pkg/broker/config/config.go @@ -59,6 +59,16 @@ type BrokerConfig struct { LeaseSeconds int `mapstructure:"lease_seconds"` // StaleJobRecoverySec is the interval between broker_recover_stale_jobs sweeps. StaleJobRecoverySec int `mapstructure:"stale_job_recovery_sec"` + // MetricsEnabled controls whether the embedded Prometheus metrics HTTP + // server (exposition endpoint + HTML dashboard) is started. + MetricsEnabled bool `mapstructure:"metrics_enabled"` + // MetricsHost is the bind address for the metrics server. + MetricsHost string `mapstructure:"metrics_host"` + // MetricsPort is the bind port for the metrics server. + MetricsPort int `mapstructure:"metrics_port"` + // QueueDepthPollSec is the interval between polls of pending job counts + // used to populate the broker_jobs_queued gauge. + QueueDepthPollSec int `mapstructure:"queue_depth_poll_sec"` } // LoggingConfig holds logging settings @@ -124,6 +134,10 @@ func setDefaults(v *viper.Viper) { v.SetDefault("broker.enable_debug", false) v.SetDefault("broker.lease_seconds", 60) v.SetDefault("broker.stale_job_recovery_sec", 30) + v.SetDefault("broker.metrics_enabled", true) + v.SetDefault("broker.metrics_host", "0.0.0.0") + v.SetDefault("broker.metrics_port", 9469) + v.SetDefault("broker.queue_depth_poll_sec", 15) // Logging defaults v.SetDefault("logging.level", "info") diff --git a/pkg/broker/config/dbedit.go b/pkg/broker/config/dbedit.go index 82ca476..8f0ff09 100644 --- a/pkg/broker/config/dbedit.go +++ b/pkg/broker/config/dbedit.go @@ -160,9 +160,10 @@ func (f *FileDoc) FindDatabase(name string) (DatabaseConfig, bool, error) { if err != nil { return DatabaseConfig{}, false, err } - for _, db := range dbs { + for i := range dbs { + db := &dbs[i] if db.Name == name { - return db, true, nil + return *db, true, nil } } return DatabaseConfig{}, false, nil diff --git a/pkg/broker/database_instance.go b/pkg/broker/database_instance.go index de1433c..d42276a 100644 --- a/pkg/broker/database_instance.go +++ b/pkg/broker/database_instance.go @@ -13,6 +13,7 @@ import ( "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/config" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/install" + "git.warky.dev/wdevs/pgsql-broker/pkg/broker/metrics" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/models" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/queue" ) @@ -36,6 +37,7 @@ type DatabaseInstance struct { shutdownMu sync.RWMutex jobsHandled int64 startTime time.Time + metrics *metrics.Metrics // sessionConn holds the pg_try_advisory_lock acquired by // registerInstance. The lock is scoped to this one physical connection, @@ -45,12 +47,16 @@ type DatabaseInstance struct { } // NewDatabaseInstance creates a new database instance -func NewDatabaseInstance(cfg *config.Config, dbCfg *config.DatabaseConfig, db adapter.DBAdapter, logger adapter.Logger, version string, parentCtx context.Context) (*DatabaseInstance, error) { +func NewDatabaseInstance(cfg *config.Config, dbCfg *config.DatabaseConfig, db adapter.DBAdapter, logger adapter.Logger, version string, parentCtx context.Context, brokerMetrics ...*metrics.Metrics) (*DatabaseInstance, error) { hostname, err := os.Hostname() if err != nil { hostname = "unknown" } + var instanceMetrics *metrics.Metrics + if len(brokerMetrics) > 0 { + instanceMetrics = brokerMetrics[0] + } instance := &DatabaseInstance{ Name: fmt.Sprintf("%s-%s", cfg.Broker.Name, dbCfg.Name), DatabaseName: dbCfg.Name, @@ -64,6 +70,7 @@ func NewDatabaseInstance(cfg *config.Config, dbCfg *config.DatabaseConfig, db ad queues: make(map[int]*queue.Queue), ctx: parentCtx, startTime: time.Now(), + metrics: instanceMetrics, } return instance, nil @@ -109,9 +116,54 @@ func (i *DatabaseInstance) Start() error { adapter.SupervisedGo(i.logger, "stale-job-recovery-routine", i.staleJobRecoveryRoutine) i.logger.Info("database instance started successfully") + if i.metrics != nil { + adapter.SupervisedGo(i.logger, "metrics-queue-depth-routine", i.queueDepthRoutine) + } return nil } +// queueDepthRoutine periodically exports pending jobs grouped by queue. +func (i *DatabaseInstance) queueDepthRoutine() { + interval := time.Duration(i.config.Broker.QueueDepthPollSec) * time.Second + if interval <= 0 { + interval = 15 * time.Second + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + i.updateQueueDepthMetrics() + for { + select { + case <-ticker.C: + i.updateQueueDepthMetrics() + case <-i.ctx.Done(): + return + } + } +} + +func (i *DatabaseInstance) updateQueueDepthMetrics() { + rows, err := i.db.Query(i.ctx, "SELECT job_queue, COUNT(*) FROM broker.broker_jobs WHERE complete_status = 0 GROUP BY job_queue") + if err != nil { + i.logger.Warn("failed to collect queue depth metrics", "error", err) + return + } + defer rows.Close() + for queueNumber := 1; queueNumber <= i.dbConfig.QueueCount; queueNumber++ { + i.metrics.SetJobsQueued(i.DatabaseName, queueNumber, 0) + } + for rows.Next() { + var queueNumber, count int + if err := rows.Scan(&queueNumber, &count); err != nil { + i.logger.Warn("failed to scan queue depth metric", "error", err) + return + } + i.metrics.SetJobsQueued(i.DatabaseName, queueNumber, float64(count)) + } + if err := rows.Err(); err != nil { + i.logger.Warn("failed to read queue depth metrics", "error", err) + } +} + // ensureSchema checks the embedded migration set against the database and, // depending on dbConfig.AutoMigrate, either applies pending migrations or // fails startup fast rather than running against a stale/missing schema. @@ -261,6 +313,8 @@ func (i *DatabaseInstance) startQueues() error { FetchSize: i.config.Broker.FetchQueryQueSize, TenantID: i.dbConfig.TenantID, LeaseSeconds: leaseSeconds, + Metrics: i.metrics, + DatabaseName: i.DatabaseName, } q := queue.New(queueCfg) @@ -271,6 +325,7 @@ func (i *DatabaseInstance) startQueues() error { i.queues[queueNum] = q i.logger.Info("queue started", "number", queueNum) } + i.metrics.SetQueueCount(i.DatabaseName, len(i.queues)) return nil } diff --git a/pkg/broker/metrics/metrics.go b/pkg/broker/metrics/metrics.go new file mode 100644 index 0000000..ddcd4db --- /dev/null +++ b/pkg/broker/metrics/metrics.go @@ -0,0 +1,130 @@ +// Package metrics defines the broker's Prometheus collectors and a small +// embedded HTTP server that exposes them, both as the standard /metrics +// exposition endpoint and as a single-page HTML dashboard. +package metrics + +import ( + "strconv" + "time" + + "github.com/prometheus/client_golang/prometheus" +) + +// Metrics holds every Prometheus collector the broker exposes. It is safe +// for concurrent use, and a nil *Metrics is safe to call methods on (all +// recording methods become no-ops), so callers that construct a Broker +// without metrics enabled don't need to special-case it. +type Metrics struct { + Registry *prometheus.Registry + + jobsCompleted *prometheus.CounterVec + jobsFailed *prometheus.CounterVec + jobsRequeued *prometheus.CounterVec + jobDuration *prometheus.HistogramVec + jobsQueued *prometheus.GaugeVec + databaseCount prometheus.Gauge + queueCount *prometheus.GaugeVec +} + +// New creates a Metrics instance with a fresh (non-global) registry, so +// multiple brokers can coexist in the same process -- e.g. in tests -- +// without colliding on prometheus.DefaultRegisterer. +func New() *Metrics { + registry := prometheus.NewRegistry() + + m := &Metrics{ + Registry: registry, + jobsCompleted: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "broker_jobs_completed_total", + Help: "Total number of jobs that completed successfully.", + }, []string{"database", "job_group", "job_name"}), + jobsFailed: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "broker_jobs_failed_total", + Help: "Total number of jobs that were dead-lettered after exhausting retries.", + }, []string{"database", "job_group", "job_name"}), + jobsRequeued: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "broker_jobs_requeued_total", + Help: "Total number of job attempts that failed and were requeued for retry.", + }, []string{"database", "job_group", "job_name"}), + jobDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "broker_job_duration_seconds", + Help: "Job execution duration in seconds, by group and name.", + Buckets: prometheus.DefBuckets, + }, []string{"database", "job_group", "job_name"}), + jobsQueued: prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Name: "broker_jobs_queued", + Help: "Current number of pending (not yet claimed) jobs, by database and queue.", + }, []string{"database", "queue"}), + databaseCount: prometheus.NewGauge(prometheus.GaugeOpts{ + Name: "broker_databases", + Help: "Number of database instances managed by this broker process.", + }), + queueCount: prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Name: "broker_queues", + Help: "Number of queues configured for a database instance.", + }, []string{"database"}), + } + + registry.MustRegister( + m.jobsCompleted, + m.jobsFailed, + m.jobsRequeued, + m.jobDuration, + m.jobsQueued, + m.databaseCount, + m.queueCount, + ) + + return m +} + +// RecordJobCompleted records a successfully completed job attempt. +func (m *Metrics) RecordJobCompleted(database, group, name string, duration time.Duration) { + if m == nil { + return + } + m.jobsCompleted.WithLabelValues(database, group, name).Inc() + m.jobDuration.WithLabelValues(database, group, name).Observe(duration.Seconds()) +} + +// RecordJobFailed records a job attempt that was dead-lettered (attempts exhausted). +func (m *Metrics) RecordJobFailed(database, group, name string, duration time.Duration) { + if m == nil { + return + } + m.jobsFailed.WithLabelValues(database, group, name).Inc() + m.jobDuration.WithLabelValues(database, group, name).Observe(duration.Seconds()) +} + +// RecordJobRequeued records a job attempt that failed but was requeued for retry. +func (m *Metrics) RecordJobRequeued(database, group, name string, duration time.Duration) { + if m == nil { + return + } + m.jobsRequeued.WithLabelValues(database, group, name).Inc() + m.jobDuration.WithLabelValues(database, group, name).Observe(duration.Seconds()) +} + +// SetJobsQueued sets the current pending job count for a database/queue pair. +func (m *Metrics) SetJobsQueued(database string, queue int, count float64) { + if m == nil { + return + } + m.jobsQueued.WithLabelValues(database, strconv.Itoa(queue)).Set(count) +} + +// SetDatabaseCount sets the number of database instances managed by this process. +func (m *Metrics) SetDatabaseCount(n int) { + if m == nil { + return + } + m.databaseCount.Set(float64(n)) +} + +// SetQueueCount sets the number of queues configured for a database instance. +func (m *Metrics) SetQueueCount(database string, n int) { + if m == nil { + return + } + m.queueCount.WithLabelValues(database).Set(float64(n)) +} diff --git a/pkg/broker/metrics/page.html b/pkg/broker/metrics/page.html new file mode 100644 index 0000000..f5aa4b3 --- /dev/null +++ b/pkg/broker/metrics/page.html @@ -0,0 +1,215 @@ + + + + + +pgsql-broker metrics + + + +

pgsql-broker metrics

+
loading…
+
+ + + + diff --git a/pkg/broker/metrics/server.go b/pkg/broker/metrics/server.go new file mode 100644 index 0000000..691937a --- /dev/null +++ b/pkg/broker/metrics/server.go @@ -0,0 +1,74 @@ +package metrics + +import ( + "context" + _ "embed" + "fmt" + "net" + "net/http" + + "github.com/prometheus/client_golang/prometheus/promhttp" + + "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" +) + +//go:embed page.html +var dashboardHTML []byte + +// Server is the embedded HTTP server exposing /metrics (standard Prometheus +// exposition format) and / (a single-page HTML dashboard that polls +// /metrics). +type Server struct { + httpServer *http.Server + listener net.Listener + logger adapter.Logger +} + +// NewServer builds a Server bound to addr (e.g. "127.0.0.1:9469"). Binding +// happens immediately so a port conflict is reported to the caller rather +// than surfacing later in a background goroutine. +func NewServer(m *Metrics, addr string, logger adapter.Logger) (*Server, error) { + listener, err := net.Listen("tcp", addr) + if err != nil { + return nil, fmt.Errorf("failed to bind metrics server to %s: %w", addr, err) + } + + mux := http.NewServeMux() + mux.Handle("/metrics", promhttp.HandlerFor(m.Registry, promhttp.HandlerOpts{})) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/" { + http.NotFound(w, r) + return + } + w.Header().Set("Content-Type", "text/html; charset=utf-8") + _, _ = w.Write(dashboardHTML) + }) + + return &Server{ + httpServer: &http.Server{Handler: mux}, + listener: listener, + logger: logger.With("component", "metrics-server"), + }, nil +} + +// Addr returns the actual bound address (useful when addr was given with a +// ":0" port). +func (s *Server) Addr() string { + return s.listener.Addr().String() +} + +// Start serves in the background. It returns immediately; Serve errors +// (other than a clean Shutdown) are logged. +func (s *Server) Start() { + s.logger.Info("metrics server listening", "addr", s.Addr()) + go func() { + if err := s.httpServer.Serve(s.listener); err != nil && err != http.ErrServerClosed { + s.logger.Error("metrics server stopped unexpectedly", "error", err) + } + }() +} + +// Stop gracefully shuts down the server. +func (s *Server) Stop(ctx context.Context) error { + return s.httpServer.Shutdown(ctx) +} diff --git a/pkg/broker/queue/queue.go b/pkg/broker/queue/queue.go index a121740..796c4fd 100644 --- a/pkg/broker/queue/queue.go +++ b/pkg/broker/queue/queue.go @@ -6,6 +6,7 @@ import ( "sync" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" + "git.warky.dev/wdevs/pgsql-broker/pkg/broker/metrics" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/worker" ) @@ -34,6 +35,8 @@ type Config struct { FetchSize int TenantID string LeaseSeconds int + Metrics *metrics.Metrics + DatabaseName string } // New creates a new queue manager @@ -73,6 +76,8 @@ func (q *Queue) Start(cfg Config) error { FetchSize: cfg.FetchSize, TenantID: cfg.TenantID, LeaseSeconds: cfg.LeaseSeconds, + Metrics: cfg.Metrics, + DatabaseName: cfg.DatabaseName, }) if err := w.Start(q.ctx); err != nil { diff --git a/pkg/broker/worker/worker.go b/pkg/broker/worker/worker.go index fa53951..7e473cb 100644 --- a/pkg/broker/worker/worker.go +++ b/pkg/broker/worker/worker.go @@ -9,6 +9,7 @@ import ( "time" "git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter" + "git.warky.dev/wdevs/pgsql-broker/pkg/broker/metrics" ) // Worker represents a single job processing worker @@ -29,6 +30,8 @@ type Worker struct { fetchSize int tenantID string leaseSeconds int + metrics *metrics.Metrics + databaseName string } // Stats holds worker statistics @@ -50,6 +53,8 @@ type Config struct { FetchSize int TenantID string LeaseSeconds int + Metrics *metrics.Metrics + DatabaseName string } // New creates a new worker @@ -72,6 +77,8 @@ func New(cfg Config) *Worker { fetchSize: cfg.FetchSize, tenantID: cfg.TenantID, leaseSeconds: leaseSeconds, + metrics: cfg.Metrics, + databaseName: cfg.DatabaseName, } } @@ -228,8 +235,19 @@ func (w *Worker) processJobs(ctx context.Context) { return // No more jobs } + jobName, jobGroup, err := w.fetchJobLabelsTx(ctx, tx, jobID) + if err != nil { + w.logger.Warn("failed to fetch job labels for metrics", "job_id", jobID, "error", err) + } + // Run the job - if err := w.runJobTx(ctx, tx, jobID, leaseToken); err != nil { + start := time.Now() + jobStatus, err := w.runJobTx(ctx, tx, jobID, leaseToken) + duration := time.Since(start) + if err == nil { + w.recordJobMetric(jobStatus, jobGroup, jobName, duration) + } + if err != nil { // Rollback on genuine infra failure if rbErr := tx.Rollback(); rbErr != nil { w.logger.Error("failed to rollback transaction", "error", rbErr) @@ -287,7 +305,7 @@ func (w *Worker) fetchNextJobTx(ctx context.Context, tx adapter.DBTransaction) ( // error (triggering a rollback of the claim) on a genuine infra failure -- // job outcomes reported via p_job_status (requeued/completed/dead-lettered) // are always committed. -func (w *Worker) runJobTx(ctx context.Context, tx adapter.DBTransaction, jobID int64, leaseToken string) error { +func (w *Worker) runJobTx(ctx context.Context, tx adapter.DBTransaction, jobID int64, leaseToken string) (int, error) { w.logger.Debug("running job", "job_id", jobID) var retval int @@ -300,15 +318,43 @@ func (w *Worker) runJobTx(ctx context.Context, tx adapter.DBTransaction, jobID i ).Scan(&retval, &errmsg, &jobStatus) if err != nil { - return fmt.Errorf("query error: %w", err) + return 0, fmt.Errorf("query error: %w", err) } if retval > 0 { - return fmt.Errorf("broker_run error: %s", errmsg) + return 0, fmt.Errorf("broker_run error: %s", errmsg) } w.logger.Debug("job finished", "job_id", jobID, "job_status", jobStatus) - return nil + return jobStatus, nil +} + +// fetchJobLabelsTx looks up the job_name/job_group of jobID for metric +// labeling. Best-effort: callers log and continue on error rather than +// failing the job over a metrics lookup. +func (w *Worker) fetchJobLabelsTx(ctx context.Context, tx adapter.DBTransaction, jobID int64) (jobName, jobGroup string, err error) { + err = tx.QueryRow(ctx, + "SELECT job_name, job_group FROM broker.broker_jobs WHERE id_broker_jobs = $1", + jobID, + ).Scan(&jobName, &jobGroup) + if err != nil { + return "", "", fmt.Errorf("query error: %w", err) + } + return jobName, jobGroup, nil +} + +// recordJobMetric routes a finished job attempt to the appropriate +// Prometheus counter/histogram based on the p_job_status broker_run +// reported (0=requeued, 2=completed, 3=dead-lettered). +func (w *Worker) recordJobMetric(jobStatus int, jobGroup, jobName string, duration time.Duration) { + switch jobStatus { + case 2: + w.metrics.RecordJobCompleted(w.databaseName, jobGroup, jobName, duration) + case 3: + w.metrics.RecordJobFailed(w.databaseName, jobGroup, jobName, duration) + case 0: + w.metrics.RecordJobRequeued(w.databaseName, jobGroup, jobName, duration) + } } // updateActivity updates the last activity timestamp