Files
EvoBGP/internal/pgmonitor/scheduler.go
T
DenozordecandCursor 8fe74c1d3b feat(jobs): add durable PG queue reclaim, slog, and richer metrics
JSON slog в ключевых пакетах; Prometheus path_group, job_audit_depth, upstream breaker; job_audit ClaimQueued/ReclaimStaleRunning + Adopt loop для HA после рестарта.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-31 12:26:18 +07:00

61 lines
1.6 KiB
Go

package pgmonitor
import (
"context"
"evobgp/internal/logging"
"fmt"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
// StartScheduler runs periodic PostgreSQL analyzer snapshots until ctx is cancelled.
func StartScheduler(ctx context.Context, pool *pgxpool.Pool) {
if pool == nil {
return
}
go func() {
t5 := time.NewTicker(5 * time.Minute)
t15 := time.NewTicker(15 * time.Minute)
defer t5.Stop()
defer t15.Stop()
s := NewService(pool)
runLight := func() {
c, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
if err := s.RefreshMetricsSnapshot(c); err != nil {
logging.Default().Info(fmt.Sprintf("pgmonitor: metrics refresh: %v", err))
}
if err := s.DetectAutovacuumLag(c); err != nil {
logging.Default().Info(fmt.Sprintf("pgmonitor: autovacuum lag: %v", err))
}
}
runHeavy := func() {
c, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
if err := s.AggregateSlowQueries(c, 30); err != nil {
logging.Default().Info(fmt.Sprintf("pgmonitor: slow queries snapshot: %v", err))
}
if err := s.EstimateTableBloat(c); err != nil {
logging.Default().Info(fmt.Sprintf("pgmonitor: bloat: %v", err))
}
if err := s.AnalyzeIndexUsage(c); err != nil {
logging.Default().Info(fmt.Sprintf("pgmonitor: index usage: %v", err))
}
}
runLight()
runHeavy()
for {
select {
case <-ctx.Done():
return
case <-t5.C:
runLight()
case <-t15.C:
runHeavy()
}
}
}()
logging.Default().Info(fmt.Sprintf("pgmonitor: scheduler started (5m light / 15m heavy)"))
}