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 }