Public source validation / validate (push) Failing after 3m8s
274 lines
7.0 KiB
Go
274 lines
7.0 KiB
Go
package alertworker
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/itworx/pulse/internal/alert"
|
|
)
|
|
|
|
const jobType = "alert-evaluation"
|
|
|
|
var (
|
|
ErrInvalidConfig = errors.New("alert evaluator configuration is invalid")
|
|
ErrLeaseLost = errors.New("alert evaluator lease was lost")
|
|
)
|
|
|
|
type RuleSource interface {
|
|
ListEnabled(context.Context, int) ([]alert.Rule, error)
|
|
}
|
|
|
|
type Evaluator interface {
|
|
Evaluate(context.Context, alert.Rule) error
|
|
}
|
|
|
|
type EvaluateFunc func(context.Context, alert.Rule) error
|
|
|
|
func (f EvaluateFunc) Evaluate(ctx context.Context, rule alert.Rule) error {
|
|
return f(ctx, rule)
|
|
}
|
|
|
|
type Config struct {
|
|
MaxConcurrent int
|
|
MaxBatch int
|
|
AttemptTimeout time.Duration
|
|
LeaseTTL time.Duration
|
|
Owner string
|
|
Now func() time.Time
|
|
}
|
|
|
|
func (c Config) validate() error {
|
|
if c.MaxConcurrent < 1 || c.MaxConcurrent > 64 || c.MaxBatch < 1 || c.MaxBatch > 100 || c.MaxConcurrent > c.MaxBatch {
|
|
return ErrInvalidConfig
|
|
}
|
|
if c.AttemptTimeout < time.Millisecond || c.AttemptTimeout > 2*time.Minute || c.LeaseTTL < c.AttemptTimeout || c.LeaseTTL > 10*time.Minute {
|
|
return ErrInvalidConfig
|
|
}
|
|
if c.Owner == "" || len(c.Owner) > 120 || c.Now == nil {
|
|
return ErrInvalidConfig
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type Lease struct {
|
|
ID string
|
|
JobKey string
|
|
ScheduledAt time.Time
|
|
Owner string
|
|
}
|
|
|
|
type LeaseStore interface {
|
|
Acquire(context.Context, string, string, time.Time, string, time.Time, time.Duration) (Lease, bool, error)
|
|
Complete(context.Context, Lease, string, string) error
|
|
}
|
|
|
|
type JobResult struct {
|
|
RuleID string
|
|
JobKey string
|
|
ScheduledAt time.Time
|
|
StartedAt time.Time
|
|
CompletedAt time.Time
|
|
Status string
|
|
ErrorCode string
|
|
}
|
|
|
|
type RunReport struct {
|
|
Scheduled int
|
|
Started int
|
|
Completed int
|
|
Skipped int
|
|
Failed int
|
|
Canceled int
|
|
Jobs []JobResult
|
|
}
|
|
|
|
type Metrics struct {
|
|
RunsStarted uint64
|
|
RunsCompleted uint64
|
|
JobsStarted uint64
|
|
JobsCompleted uint64
|
|
JobsFailed uint64
|
|
JobsSkipped uint64
|
|
LastRunDuration time.Duration
|
|
}
|
|
|
|
type Worker struct {
|
|
Source RuleSource
|
|
Store LeaseStore
|
|
Evaluator Evaluator
|
|
Config Config
|
|
mu sync.Mutex
|
|
metrics Metrics
|
|
}
|
|
|
|
func New(source RuleSource, store LeaseStore, evaluator Evaluator, config Config) (Worker, error) {
|
|
if source == nil || store == nil || evaluator == nil {
|
|
return Worker{}, ErrInvalidConfig
|
|
}
|
|
if config.Owner == "" {
|
|
config.Owner = alert.NewID()
|
|
}
|
|
if config.Now == nil {
|
|
config.Now = time.Now
|
|
}
|
|
if err := config.validate(); err != nil {
|
|
return Worker{}, err
|
|
}
|
|
return Worker{Source: source, Store: store, Evaluator: evaluator, Config: config}, nil
|
|
}
|
|
|
|
func (w *Worker) RunOnce(ctx context.Context) (RunReport, error) {
|
|
if err := ctx.Err(); err != nil {
|
|
return RunReport{}, err
|
|
}
|
|
start := w.Config.Now().UTC()
|
|
w.mu.Lock()
|
|
w.metrics.RunsStarted++
|
|
w.mu.Unlock()
|
|
rules, err := w.Source.ListEnabled(ctx, w.Config.MaxBatch)
|
|
if err != nil {
|
|
return RunReport{}, fmt.Errorf("list enabled alert rules: %w", err)
|
|
}
|
|
report := RunReport{Scheduled: len(rules), Jobs: make([]JobResult, 0, len(rules))}
|
|
type job struct {
|
|
rule alert.Rule
|
|
lease Lease
|
|
}
|
|
jobs := make([]job, 0, len(rules))
|
|
for _, rule := range rules {
|
|
interval := time.Duration(rule.EvaluationIntervalSeconds) * time.Second
|
|
scheduledAt := start.Truncate(interval)
|
|
key := rule.ID
|
|
lease, acquired, err := w.Store.Acquire(ctx, jobType, key, scheduledAt, w.Config.Owner, start, w.Config.LeaseTTL)
|
|
if err != nil {
|
|
return report, fmt.Errorf("acquire alert evaluation lease: %w", err)
|
|
}
|
|
if !acquired {
|
|
report.Skipped++
|
|
w.mu.Lock()
|
|
w.metrics.JobsSkipped++
|
|
w.mu.Unlock()
|
|
report.Jobs = append(report.Jobs, JobResult{RuleID: rule.ID, JobKey: key, ScheduledAt: scheduledAt, Status: "skipped"})
|
|
continue
|
|
}
|
|
jobs = append(jobs, job{rule: rule, lease: lease})
|
|
}
|
|
if len(jobs) == 0 {
|
|
w.finishRun(start)
|
|
return report, nil
|
|
}
|
|
sem := make(chan struct{}, w.Config.MaxConcurrent)
|
|
var wait sync.WaitGroup
|
|
var reportMu sync.Mutex
|
|
for _, item := range jobs {
|
|
select {
|
|
case <-ctx.Done():
|
|
report.Canceled++
|
|
report.Jobs = append(report.Jobs, JobResult{RuleID: item.rule.ID, JobKey: item.lease.JobKey, ScheduledAt: item.lease.ScheduledAt, Status: "canceled", ErrorCode: "shutdown"})
|
|
_ = w.completeLease(item.lease, "canceled", "shutdown")
|
|
case sem <- struct{}{}:
|
|
wait.Add(1)
|
|
report.Started++
|
|
w.mu.Lock()
|
|
w.metrics.JobsStarted++
|
|
w.mu.Unlock()
|
|
go func(item job) {
|
|
defer wait.Done()
|
|
defer func() { <-sem }()
|
|
jobResult := w.evaluate(ctx, item.rule, item.lease)
|
|
reportMu.Lock()
|
|
report.Jobs = append(report.Jobs, jobResult)
|
|
switch jobResult.Status {
|
|
case "completed":
|
|
report.Completed++
|
|
case "failed":
|
|
report.Failed++
|
|
case "canceled":
|
|
report.Canceled++
|
|
}
|
|
reportMu.Unlock()
|
|
}(item)
|
|
}
|
|
}
|
|
wait.Wait()
|
|
sort.Slice(report.Jobs, func(i, j int) bool {
|
|
if report.Jobs[i].ScheduledAt.Equal(report.Jobs[j].ScheduledAt) {
|
|
return report.Jobs[i].RuleID < report.Jobs[j].RuleID
|
|
}
|
|
return report.Jobs[i].ScheduledAt.Before(report.Jobs[j].ScheduledAt)
|
|
})
|
|
w.finishRun(start)
|
|
return report, nil
|
|
}
|
|
|
|
func (w *Worker) RunLoop(ctx context.Context, interval time.Duration) error {
|
|
if interval < time.Second || interval > time.Hour {
|
|
return ErrInvalidConfig
|
|
}
|
|
for {
|
|
_, err := w.RunOnce(ctx)
|
|
if err != nil && !errors.Is(err, context.Canceled) {
|
|
return err
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case <-time.After(interval):
|
|
}
|
|
}
|
|
}
|
|
|
|
func (w *Worker) evaluate(ctx context.Context, rule alert.Rule, lease Lease) JobResult {
|
|
started := w.Config.Now().UTC()
|
|
jobResult := JobResult{RuleID: rule.ID, JobKey: lease.JobKey, ScheduledAt: lease.ScheduledAt, StartedAt: started}
|
|
attemptCtx, cancel := context.WithTimeout(ctx, w.Config.AttemptTimeout)
|
|
err := w.Evaluator.Evaluate(attemptCtx, rule)
|
|
cancel()
|
|
status, errorCode := "completed", ""
|
|
if err != nil {
|
|
status = "failed"
|
|
errorCode = "evaluation_failed"
|
|
if errors.Is(err, context.DeadlineExceeded) || errors.Is(attemptCtx.Err(), context.DeadlineExceeded) {
|
|
errorCode = "timeout"
|
|
} else if errors.Is(err, context.Canceled) || errors.Is(ctx.Err(), context.Canceled) {
|
|
status, errorCode = "canceled", "shutdown"
|
|
}
|
|
}
|
|
jobResult.CompletedAt = w.Config.Now().UTC()
|
|
jobResult.Status = status
|
|
jobResult.ErrorCode = errorCode
|
|
_ = w.completeLease(lease, status, errorCode)
|
|
w.mu.Lock()
|
|
switch status {
|
|
case "completed":
|
|
w.metrics.JobsCompleted++
|
|
case "failed", "canceled":
|
|
w.metrics.JobsFailed++
|
|
}
|
|
w.mu.Unlock()
|
|
return jobResult
|
|
}
|
|
|
|
func (w *Worker) completeLease(lease Lease, status, errorCode string) error {
|
|
ctx, cancel := context.WithTimeout(context.WithoutCancel(context.Background()), 2*time.Second)
|
|
defer cancel()
|
|
return w.Store.Complete(ctx, lease, status, errorCode)
|
|
}
|
|
|
|
func (w *Worker) finishRun(start time.Time) {
|
|
w.mu.Lock()
|
|
w.metrics.RunsCompleted++
|
|
w.metrics.LastRunDuration = w.Config.Now().UTC().Sub(start)
|
|
w.mu.Unlock()
|
|
}
|
|
|
|
func (w *Worker) Metrics() Metrics {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
return w.metrics
|
|
}
|