Files
atlantis/server/scheduled/executor_service.go
Ken Kaizu 3954955e13 fix(deps): update module github.com/petergtz/pegomock/v3 to v4 (#3534)
* fix(deps): update module github.com/petergtz/pegomock/v3 to v4 in go.mod

* remove pegomock generate m option, which is not support after v4

* make regen-mocks

* replace pegomock v4 primitive eq/matchers

* convert pegomock v4 Eq/Any matchers

* remove custom models.Repo matcher

* pegomock v4 cannot use result method args

ref https://github.com/petergtz/pegomock/issues/123

---------

Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com>
2023-06-20 15:05:23 -04:00

109 lines
2.1 KiB
Go

package scheduled
import (
"context"
"os"
"os/signal"
"sync"
"syscall"
"time"
"github.com/runatlantis/atlantis/server/logging"
tally "github.com/uber-go/tally/v4"
)
type ExecutorService struct {
log logging.SimpleLogging
// jobs
jobs []JobDefinition
}
func NewExecutorService(
statsScope tally.Scope,
log logging.SimpleLogging,
) *ExecutorService {
scheduledScope := statsScope.SubScope("scheduled")
runtimeStatsPublisher := NewRuntimeStats(scheduledScope)
runtimeStatsPublisherJob := JobDefinition{
Job: runtimeStatsPublisher,
Period: 10 * time.Second,
}
return &ExecutorService{
log: log,
jobs: []JobDefinition{runtimeStatsPublisherJob},
}
}
func (s *ExecutorService) AddJob(jd JobDefinition) {
s.jobs = append(s.jobs, jd)
}
type JobDefinition struct {
Job Job
Period time.Duration
}
func (s *ExecutorService) Run() {
s.log.Info("Scheduled Executor Service started")
ctx, cancel := context.WithCancel(context.Background())
var wg sync.WaitGroup
for _, jd := range s.jobs {
s.runScheduledJob(ctx, &wg, jd)
}
interrupt := make(chan os.Signal, 1)
// Stop on SIGINTs and SIGTERMs.
signal.Notify(interrupt, os.Interrupt, syscall.SIGTERM)
<-interrupt
s.log.Warn("Received interrupt. Attempting to Shut down scheduled executor service")
cancel()
wg.Wait()
s.log.Warn("All jobs completed, exiting.")
}
func (s *ExecutorService) runScheduledJob(ctx context.Context, wg *sync.WaitGroup, jd JobDefinition) {
ticker := time.NewTicker(jd.Period)
wg.Add(1)
go func() {
defer wg.Done()
defer ticker.Stop()
// Ensure we recover from any panics to keep the jobs isolated.
// Keep the recovery outside the select to ensure that we don't infinitely panic.
defer func() {
if r := recover(); r != nil {
s.log.Err("Recovered from panic: %v", r)
}
}()
for {
select {
case <-ctx.Done():
s.log.Warn("Received interrupt, cancelling job")
return
case <-ticker.C:
jd.Job.Run()
}
}
}()
}
//go:generate pegomock generate --package mocks -o mocks/mock_executor_service_job.go Job
type Job interface {
Run()
}