60 lines
1.3 KiB
Go
60 lines
1.3 KiB
Go
package scheduler
|
|
|
|
import (
|
|
"context"
|
|
"log/slog"
|
|
"time"
|
|
)
|
|
|
|
type Job struct {
|
|
Name string
|
|
Interval time.Duration
|
|
Fn func(ctx context.Context) error
|
|
}
|
|
|
|
type Scheduler struct {
|
|
jobs []Job
|
|
logger *slog.Logger
|
|
}
|
|
|
|
func New(logger *slog.Logger) *Scheduler {
|
|
return &Scheduler{logger: logger}
|
|
}
|
|
|
|
func (s *Scheduler) Add(job Job) {
|
|
s.jobs = append(s.jobs, job)
|
|
}
|
|
|
|
// Start lance tous les jobs en arrière-plan. Chaque job est exécuté immédiatement
|
|
// puis répété selon son intervalle. S'arrête proprement à l'annulation du contexte.
|
|
func (s *Scheduler) Start(ctx context.Context) {
|
|
for _, job := range s.jobs {
|
|
go s.run(ctx, job)
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) run(ctx context.Context, job Job) {
|
|
s.logger.Info("scheduler: starting job", "name", job.Name, "interval", job.Interval)
|
|
|
|
s.execute(ctx, job)
|
|
|
|
ticker := time.NewTicker(job.Interval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
s.logger.Info("scheduler: job stopped", "name", job.Name)
|
|
return
|
|
case <-ticker.C:
|
|
s.execute(ctx, job)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Scheduler) execute(ctx context.Context, job Job) {
|
|
s.logger.Info("scheduler: running job", "name", job.Name)
|
|
if err := job.Fn(ctx); err != nil {
|
|
s.logger.Error("scheduler: job failed", "name", job.Name, "error", err)
|
|
}
|
|
}
|