scheduler
packageAPI reference for the scheduler
package.
Imports
(14)Job
Job defines a scheduled task with a cron expression and handler.
type Job struct
Fields
| Name | Type | Description |
|---|---|---|
| Name | string | |
| Cron | string | |
| Handler | func(ctx context.Context) error |
Scheduler
Scheduler runs jobs according to cron expressions.
type Scheduler struct
Methods
SetLogger replaces the scheduler logging callback.
Parameters
func (*Scheduler) SetLogger(log func(msg string))
{
s.mu.Lock()
defer s.mu.Unlock()
s.logger = log
}
Parameters
func (*Scheduler) log(msg string)
{
s.mu.RLock()
logger := s.logger
s.mu.RUnlock()
if logger != nil {
func() {
defer func() {
_ = recover()
}()
logger(msg)
}()
}
}
Parameters
func (*Scheduler) metric(name string, dur time.Duration, err error)
{
s.mu.RLock()
metrics := s.metrics
s.mu.RUnlock()
if metrics != nil {
func() {
defer func() {
_ = recover()
}()
metrics(name, dur, err)
}()
}
}
Parameters
Returns
func (*Scheduler) Register(job Job) error
{
s.mu.Lock()
defer s.mu.Unlock()
if s.storeErr != nil {
return s.storeErr
}
if job.Name == "" {
return errors.New("scheduler: job name cannot be empty")
}
if job.Handler == nil {
return fmt.Errorf("scheduler: job %q handler cannot be nil", job.Name)
}
if _, err := parseCron(job.Cron); err != nil {
return err
}
for _, existing := range s.jobs {
if existing.job.Name == job.Name {
return fmt.Errorf("scheduler: job %q is already registered", job.Name)
}
}
var last time.Time
if s.store != nil {
rec, err := s.store.Load(job.Name)
switch {
case err == nil && !rec.LastRun.IsZero():
last = rec.LastRun
case err != nil && !errors.Is(err, os.ErrNotExist):
return fmt.Errorf("scheduler: load job %q state: %w", job.Name, err)
}
}
s.jobs = append(s.jobs, scheduledJob{job: job, last: last})
return nil
}
Parameters
Returns
func (*Scheduler) Start(ctx context.Context) error
{
s.mu.Lock()
if s.storeErr != nil {
s.mu.Unlock()
return s.storeErr
}
if s.running {
s.mu.Unlock()
return fmt.Errorf("scheduler: already running")
}
s.taskMu.Lock()
if s.stopping && s.cancel != nil {
s.taskMu.Unlock()
s.mu.Unlock()
return fmt.Errorf("scheduler: shutdown is still in progress")
}
s.stopping = false
s.taskMu.Unlock()
s.running = true
s.runErr = nil
ctx, s.cancel = context.WithCancel(ctx)
s.runCtx = ctx
s.done = make(chan struct{})
s.failure = make(chan error, 1)
done := s.done
s.mu.Unlock()
s.log("scheduler: started")
go func() {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
defer close(done)
for {
select {
case <-ctx.Done():
s.mu.Lock()
s.running = false
s.mu.Unlock()
s.log("scheduler: stopped")
return
case now := <-ticker.C:
s.runDue(ctx, now)
}
}
}()
return nil
}
Parameters
func (*Scheduler) runDue(ctx context.Context, now time.Time)
{
s.mu.RLock()
jobs := make([]scheduledJob, len(s.jobs))
copy(jobs, s.jobs)
s.mu.RUnlock()
for i, sj := range jobs {
if isDue(sj.job.Cron, now, sj.last) {
launched := s.launchTask(func() {
start := time.Now()
err := callJob(ctx, sj.job.Handler)
dur := time.Since(start)
s.metric(sj.job.Name, dur, err)
if err != nil {
s.log(fmt.Sprintf("scheduler: job %s error: %v", sj.job.Name, err))
}
if s.store != nil {
status := "ok"
if err != nil {
status = "error"
}
if saveErr := s.store.Save(&JobRecord{
Name: sj.job.Name,
Cron: sj.job.Cron,
LastRun: start,
LastStatus: status,
LastLatency: dur.String(),
}); saveErr != nil {
s.fail(fmt.Errorf("scheduler: persist job %s state: %w", sj.job.Name, saveErr))
}
}
})
if launched {
s.mu.Lock()
s.jobs[i].last = now
s.mu.Unlock()
}
}
}
}
Parameters
Returns
func (*Scheduler) Stop(ctx context.Context) error
{
s.mu.Lock()
cancel := s.cancel
done := s.done
if cancel == nil {
s.mu.Unlock()
return nil
}
s.mu.Unlock()
s.taskMu.Lock()
s.stopping = true
s.taskMu.Unlock()
cancel()
if done != nil {
select {
case <-done:
case <-ctx.Done():
return ctx.Err()
}
}
tasksDone := make(chan struct{})
go func() {
s.tasks.Wait()
close(tasksDone)
}()
select {
case <-tasksDone:
case <-ctx.Done():
return ctx.Err()
}
s.mu.Lock()
s.cancel = nil
s.runCtx = nil
s.done = nil
s.mu.Unlock()
return nil
}
Parameters
func (*Scheduler) Enqueue(fn func(ctx context.Context) error)
{
if fn == nil {
s.fail(errors.New("scheduler: enqueue handler cannot be nil"))
return
}
s.log("scheduler: enqueued fire-and-forget job")
ctx := s.taskContext()
s.launchTask(func() {
start := time.Now()
err := callJob(ctx, fn)
dur := time.Since(start)
s.metric("enqueue", dur, err)
if err != nil {
s.log(fmt.Sprintf("scheduler: enqueue error: %v", err))
}
})
}
Parameters
func (*Scheduler) ScheduleAfter(d time.Duration, fn func(ctx context.Context) error)
{
if fn == nil {
s.fail(errors.New("scheduler: schedule-after handler cannot be nil"))
return
}
s.log(fmt.Sprintf("scheduler: scheduled job after %v", d))
ctx := s.taskContext()
s.launchTask(func() {
timer := time.NewTimer(d)
defer timer.Stop()
select {
case <-timer.C:
case <-ctx.Done():
return
}
start := time.Now()
err := callJob(ctx, fn)
dur := time.Since(start)
s.metric("schedule-after", dur, err)
if err != nil {
s.log(fmt.Sprintf("scheduler: schedule-after error: %v", err))
}
})
}
Returns
func (*Scheduler) taskContext() context.Context
{
s.mu.RLock()
defer s.mu.RUnlock()
if s.runCtx != nil {
return s.runCtx
}
return context.Background()
}
Parameters
Returns
func (*Scheduler) launchTask(task func()) bool
{
s.taskMu.Lock()
if s.stopping {
s.taskMu.Unlock()
return false
}
s.tasks.Add(1)
s.taskMu.Unlock()
go func() {
defer s.tasks.Done()
defer func() {
if recovered := recover(); recovered != nil {
s.fail(fmt.Errorf("scheduler: async task panic: %v", recovered))
}
}()
task()
}()
return true
}
Parameters
func (*Scheduler) fail(err error)
{
if err == nil {
return
}
s.mu.Lock()
s.runErr = errors.Join(s.runErr, err)
cancel := s.cancel
failure := s.failure
s.mu.Unlock()
if failure != nil {
select {
case failure <- err:
default:
}
}
s.log(err.Error())
if cancel != nil {
cancel()
}
}
Err returns runtime failures that stopped the scheduler.
Returns
func (*Scheduler) Err() error
{
s.mu.RLock()
defer s.mu.RUnlock()
return s.runErr
}
Completion reports a runtime failure that stops the scheduler.
Returns
func (*Scheduler) Completion() <-chan error
{
s.mu.RLock()
defer s.mu.RUnlock()
return s.failure
}
Fields
| Name | Type | Description |
|---|---|---|
| jobs | []scheduledJob | |
| mu | sync.RWMutex | |
| running | bool | |
| cancel | context.CancelFunc | |
| logger | func(msg string) | |
| metrics | func(name string, dur time.Duration, err error) | |
| store | *JobStore | |
| storeErr | error | |
| runErr | error | |
| runCtx | context.Context | |
| done | chan struct{} | |
| failure | chan error | |
| taskMu | sync.Mutex | |
| tasks | sync.WaitGroup | |
| stopping | bool |
scheduledJob
type scheduledJob struct
Uses
Option
Option configures a Scheduler.
type Option options.Option[Scheduler]
New
New creates a new Scheduler with the given options.
Parameters
Returns
func New(opts ...Option) *Scheduler
{
s := &Scheduler{}
for _, opt := range opts {
opt(s)
}
return s
}
WithLogger
WithLogger sets a logging callback for the scheduler.
Parameters
Returns
func WithLogger(log func(msg string)) Option
{
return func(s *Scheduler) { s.logger = log }
}
Uses
WithMetrics
WithMetrics sets a metrics callback for job executions.
Parameters
Returns
func WithMetrics(m func(name string, dur time.Duration, err error)) Option
{
return func(s *Scheduler) { s.metrics = m }
}
Uses
WithStore
WithStore enables persistent job state storage to the given directory.
Parameters
Returns
func WithStore(dir string) Option
{
return func(s *Scheduler) {
store, err := NewJobStore(dir)
if err != nil {
s.storeErr = err
return
}
s.store = store
}
}
Uses
callJob
Parameters
Returns
func callJob(ctx context.Context, handler func(context.Context) error) (err error)
{
defer func() {
if recovered := recover(); recovered != nil {
err = fmt.Errorf("scheduler: job panic: %v", recovered)
}
}()
return handler(ctx)
}
JobRecord
JobRecord holds persisted job execution state.
type JobRecord struct
Fields
| Name | Type | Description |
|---|---|---|
| Name | string | json:"name" |
| Cron | string | json:"cron" |
| LastRun | time.Time | json:"last_run" |
| LastStatus | string | json:"last_status,omitempty" |
| LastLatency | string | json:"last_latency,omitempty" |
JobStore
JobStore persists job state to disk as JSON.
type JobStore struct
Methods
Parameters
Returns
func (*JobStore) Save(rec *JobRecord) error
{
js.mu.Lock()
defer js.mu.Unlock()
if rec == nil {
return errors.New("scheduler store: record cannot be nil")
}
if err := validateJobName(rec.Name); err != nil {
return err
}
data, err := json.MarshalIndent(rec, "", " ")
if err != nil {
return fmt.Errorf("scheduler store: marshal: %w", err)
}
if len(data) > maxJobRecordSize {
return errors.New("scheduler store: record exceeds size limit")
}
path := filepath.Join(js.dir, rec.Name+".json")
tmp, err := os.CreateTemp(js.dir, ".job-*.tmp")
if err != nil {
return fmt.Errorf("scheduler store: create temp: %w", err)
}
tmpName := tmp.Name()
defer os.Remove(tmpName)
if err := tmp.Chmod(0600); err != nil {
tmp.Close()
return fmt.Errorf("scheduler store: secure temp: %w", err)
}
if _, err := tmp.Write(data); err != nil {
tmp.Close()
return fmt.Errorf("scheduler store: write: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("scheduler store: close: %w", err)
}
if err := os.Rename(tmpName, path); err != nil {
return fmt.Errorf("scheduler store: rename: %w", err)
}
return nil
}
Parameters
Returns
func (*JobStore) Load(name string) (*JobRecord, error)
{
js.mu.RLock()
defer js.mu.RUnlock()
if err := validateJobName(name); err != nil {
return nil, err
}
data, err := readJobFile(js.dir, name+".json")
if err != nil {
return nil, fmt.Errorf("scheduler store: read %s: %w", name, err)
}
var rec JobRecord
if err := json.Unmarshal(data, &rec); err != nil {
return nil, fmt.Errorf("scheduler store: unmarshal %s: %w", name, err)
}
return &rec, nil
}
Returns
func (*JobStore) List() ([]*JobRecord, error)
{
js.mu.RLock()
defer js.mu.RUnlock()
entries, err := os.ReadDir(js.dir)
if err != nil {
return nil, fmt.Errorf("scheduler store: readdir: %w", err)
}
var result []*JobRecord
for _, entry := range entries {
if filepath.Ext(entry.Name()) != ".json" || entry.Type()&os.ModeSymlink != 0 {
continue
}
data, err := readJobFile(js.dir, entry.Name())
if err != nil {
continue
}
var rec JobRecord
if err := json.Unmarshal(data, &rec); err != nil {
continue
}
result = append(result, &rec)
}
return result, nil
}
Fields
| Name | Type | Description |
|---|---|---|
| dir | string | |
| mu | sync.RWMutex |
NewJobStore
NewJobStore creates a JobStore that writes to the given directory.
Parameters
Returns
func NewJobStore(dir string) (*JobStore, error)
{
if err := os.MkdirAll(dir, 0700); err != nil {
return nil, fmt.Errorf("scheduler store: cannot create dir %s: %w", dir, err)
}
if err := os.Chmod(dir, 0700); err != nil {
return nil, fmt.Errorf("scheduler store: cannot secure dir %s: %w", dir, err)
}
return &JobStore{dir: dir}, nil
}
validateJobName
Parameters
Returns
func validateJobName(name string) error
{
if name == "" || name == "." || name == ".." ||
strings.ContainsAny(name, `/\`) {
return fmt.Errorf("scheduler store: invalid job name %q", name)
}
for _, character := range name {
if character == 0 {
return fmt.Errorf("scheduler store: invalid job name %q", name)
}
}
return nil
}
readJobFile
Parameters
Returns
func readJobFile(dir, name string) ([]byte, error)
{
root, err := os.OpenRoot(dir)
if err != nil {
return nil, err
}
defer root.Close()
info, err := root.Lstat(name)
if err != nil {
return nil, err
}
if info.Mode()&os.ModeSymlink != 0 {
return nil, errors.New("scheduler store: refusing to read symbolic link")
}
file, err := root.Open(name)
if err != nil {
return nil, err
}
defer file.Close()
data, err := io.ReadAll(io.LimitReader(file, maxJobRecordSize+1))
if err != nil {
return nil, err
}
if len(data) > maxJobRecordSize {
return nil, errors.New("scheduler store: record exceeds size limit")
}
return data, nil
}
isDue
func isDue(cronExpr string, now, last time.Time) bool
{
fields, err := parseCron(cronExpr)
if err != nil {
return false
}
if last.IsZero() {
return cronMatches(fields, now)
}
next, ok := nextRun(fields, last)
return ok && !next.After(now)
}
cronFields
type cronFields struct
Fields
| Name | Type | Description |
|---|---|---|
| minute | int | |
| hour | int | |
| dom | int | |
| month | int | |
| dow | int | |
| wildMin | bool | |
| wildHour | bool | |
| wildDom | bool | |
| wildMonth | bool | |
| wildDow | bool |
parseCron
Parameters
Returns
func parseCron(expr string) (cronFields, error)
{
parts := splitFields(expr)
if len(parts) != 5 {
return cronFields{}, fmt.Errorf("scheduler: invalid cron expression: %s", expr)
}
f := cronFields{}
var err error
if f.minute, f.wildMin, err = parseField(parts[0], 0, 59); err != nil {
return cronFields{}, err
}
if f.hour, f.wildHour, err = parseField(parts[1], 0, 23); err != nil {
return cronFields{}, err
}
if f.dom, f.wildDom, err = parseField(parts[2], 1, 31); err != nil {
return cronFields{}, err
}
if f.month, f.wildMonth, err = parseField(parts[3], 1, 12); err != nil {
return cronFields{}, err
}
if f.dow, f.wildDow, err = parseField(parts[4], 0, 6); err != nil {
return cronFields{}, err
}
return f, nil
}
splitFields
Parameters
Returns
func splitFields(expr string) []string
{
var fields []string
current := ""
for _, ch := range expr {
if ch == ' ' || ch == '\t' {
if current != "" {
fields = append(fields, current)
current = ""
}
} else {
current += string(ch)
}
}
if current != "" {
fields = append(fields, current)
}
return fields
}
parseField
Parameters
Returns
func parseField(s string, min, max int) (int, bool, error)
{
if s == "*" {
return 0, true, nil
}
value, err := strconv.Atoi(s)
if err != nil || value < min || value > max {
return 0, false, fmt.Errorf("scheduler: invalid cron field %q", s)
}
return value, false, nil
}
nextRun
Parameters
Returns
func nextRun(fields cronFields, last time.Time) (time.Time, bool)
{
next := last.Truncate(time.Minute).Add(time.Minute)
const fiveYearsInMinutes = 5 * 366 * 24 * 60
for i := 0; i < fiveYearsInMinutes; i++ {
if cronMatches(fields, next) {
return next, true
}
next = next.Add(time.Minute)
}
return time.Time{}, false
}
cronMatches
Parameters
Returns
func cronMatches(fields cronFields, value time.Time) bool
{
return (fields.wildMin || value.Minute() == fields.minute) &&
(fields.wildHour || value.Hour() == fields.hour) &&
(fields.wildDom || value.Day() == fields.dom) &&
(fields.wildMonth || int(value.Month()) == fields.month) &&
(fields.wildDow || int(value.Weekday()) == fields.dow)
}
TestScheduler_RegisterAndFire
Parameters
func TestScheduler_RegisterAndFire(t *testing.T)
{
s := New()
var mu sync.Mutex
fired := false
if err := s.Register(Job{
Name: "test",
Cron: "* * * * *",
Handler: func(ctx context.Context) error { mu.Lock(); fired = true; mu.Unlock(); return nil },
}); err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
if err := s.Start(ctx); err != nil {
t.Fatalf("Start failed: %v", err)
}
time.Sleep(1500 * time.Millisecond)
s.mu.RLock()
running := s.running
s.mu.RUnlock()
if !running {
t.Error("scheduler should be running")
}
mu.Lock()
wasFired := fired
mu.Unlock()
if !wasFired {
t.Error("job should have fired within timeout")
}
}
TestScheduler_DoubleStart
Parameters
func TestScheduler_DoubleStart(t *testing.T)
{
s := New()
ctx := context.Background()
if err := s.Start(ctx); err != nil {
t.Fatalf("First Start failed: %v", err)
}
if err := s.Start(ctx); err == nil {
t.Error("double start should return an error")
}
}
TestScheduler_RejectsInvalidJobs
Parameters
func TestScheduler_RejectsInvalidJobs(t *testing.T)
{
tests := []Job{
{Name: "", Cron: "* * * * *", Handler: func(context.Context) error { return nil }},
{Name: "nil", Cron: "* * * * *"},
{Name: "bad", Cron: "bogus * * * *", Handler: func(context.Context) error { return nil }},
{Name: "range", Cron: "60 * * * *", Handler: func(context.Context) error { return nil }},
}
for _, job := range tests {
scheduler := New()
if err := scheduler.Register(job); err == nil {
t.Fatalf("Register() accepted job %#v", job)
}
}
scheduler := New()
job := Job{Name: "duplicate", Cron: "* * * * *", Handler: func(context.Context) error { return nil }}
if err := scheduler.Register(job); err != nil {
t.Fatal(err)
}
if err := scheduler.Register(job); err == nil {
t.Fatal("Register() accepted a duplicate name")
}
}
TestScheduler_DoesNotRunFutureCronImmediately
Parameters
func TestScheduler_DoesNotRunFutureCronImmediately(t *testing.T)
{
now := time.Date(2026, time.July, 26, 12, 30, 0, 0, time.UTC)
if isDue("0 0 * * *", now, time.Time{}) {
t.Fatal("future cron was due only because it had no previous run")
}
}
TestScheduler_StopNotRunning
Parameters
func TestScheduler_StopNotRunning(t *testing.T)
{
s := New()
ctx := context.Background()
if err := s.Stop(ctx); err != nil {
t.Fatalf("Stop() should be idempotent: %v", err)
}
}
TestScheduler_StopWaitsForInflightTask
Parameters
func TestScheduler_StopWaitsForInflightTask(t *testing.T)
{
scheduler := New()
if err := scheduler.Start(context.Background()); err != nil {
t.Fatal(err)
}
started := make(chan struct{})
release := make(chan struct{})
scheduler.Enqueue(func(context.Context) error {
close(started)
<-release
return nil
})
<-started
stopped := make(chan error, 1)
go func() {
stopped <- scheduler.Stop(context.Background())
}()
select {
case err := <-stopped:
t.Fatalf("Stop() returned before the task completed: %v", err)
case <-time.After(20 * time.Millisecond):
}
close(release)
if err := <-stopped; err != nil {
t.Fatal(err)
}
}
TestScheduler_RejectsRestartWhileTimedOutStopIsPending
Parameters
func TestScheduler_RejectsRestartWhileTimedOutStopIsPending(t *testing.T)
{
scheduler := New()
if err := scheduler.Start(context.Background()); err != nil {
t.Fatal(err)
}
started := make(chan struct{})
release := make(chan struct{})
scheduler.Enqueue(func(context.Context) error {
close(started)
<-release
return nil
})
<-started
stopCtx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond)
defer cancel()
if err := scheduler.Stop(stopCtx); err == nil {
t.Fatal("Stop() ignored its deadline")
}
if err := scheduler.Start(context.Background()); err == nil {
t.Fatal("Start() restarted while the previous Stop was incomplete")
}
close(release)
if err := scheduler.Stop(context.Background()); err != nil {
t.Fatal(err)
}
if err := scheduler.Start(context.Background()); err != nil {
t.Fatalf("Start() after completed Stop error = %v", err)
}
if err := scheduler.Stop(context.Background()); err != nil {
t.Fatal(err)
}
}
TestScheduler_Enqueue
Parameters
func TestScheduler_Enqueue(t *testing.T)
{
s := New()
var mu sync.Mutex
called := false
s.Enqueue(func(ctx context.Context) error {
mu.Lock()
called = true
mu.Unlock()
return nil
})
time.Sleep(100 * time.Millisecond)
mu.Lock()
wasCalled := called
mu.Unlock()
if !wasCalled {
t.Error("enqueued function should have been called")
}
}
TestScheduler_ScheduleAfter
Parameters
func TestScheduler_ScheduleAfter(t *testing.T)
{
s := New()
var mu sync.Mutex
called := false
s.ScheduleAfter(50*time.Millisecond, func(ctx context.Context) error {
mu.Lock()
called = true
mu.Unlock()
return nil
})
time.Sleep(200 * time.Millisecond)
mu.Lock()
wasCalled := called
mu.Unlock()
if !wasCalled {
t.Error("scheduled function should have been called after delay")
}
}
TestSchedulerRecoversAsyncHandlerPanics
Parameters
func TestSchedulerRecoversAsyncHandlerPanics(t *testing.T)
{
messages := make(chan string, 1)
scheduler := New(WithLogger(func(message string) {
if strings.Contains(message, "panic") {
messages <- message
}
}))
scheduler.Enqueue(func(context.Context) error {
panic("job failed")
})
select {
case message := <-messages:
if !strings.Contains(message, "job failed") {
t.Fatalf("panic log = %q", message)
}
case <-time.After(time.Second):
t.Fatal("scheduler did not report recovered panic")
}
}
TestJobStore_SaveAndLoad
Parameters
func TestJobStore_SaveAndLoad(t *testing.T)
{
dir := t.TempDir()
store, err := NewJobStore(dir)
if err != nil {
t.Fatalf("NewJobStore: %v", err)
}
rec := &JobRecord{
Name: "cleanup",
Cron: "0 3 * * *",
LastRun: time.Date(2026, 1, 15, 3, 0, 0, 0, time.UTC),
LastStatus: "ok",
LastLatency: "12ms",
}
if err := store.Save(rec); err != nil {
t.Fatalf("Save: %v", err)
}
loaded, err := store.Load("cleanup")
if err != nil {
t.Fatalf("Load: %v", err)
}
if loaded.Name != "cleanup" {
t.Errorf("Name = %q, want %q", loaded.Name, "cleanup")
}
if loaded.Cron != "0 3 * * *" {
t.Errorf("Cron = %q, want %q", loaded.Cron, "0 3 * * *")
}
if !loaded.LastRun.Equal(rec.LastRun) {
t.Errorf("LastRun = %v, want %v", loaded.LastRun, rec.LastRun)
}
}
TestJobStore_List
Parameters
func TestJobStore_List(t *testing.T)
{
dir := t.TempDir()
store, _ := NewJobStore(dir)
store.Save(&JobRecord{Name: "job1", Cron: "* * * * *"})
store.Save(&JobRecord{Name: "job2", Cron: "0 * * * *"})
list, err := store.List()
if err != nil {
t.Fatalf("List: %v", err)
}
if len(list) != 2 {
t.Errorf("got %d jobs, want 2", len(list))
}
}
TestJobStoreRejectsTraversalAndSymlinks
Parameters
func TestJobStoreRejectsTraversalAndSymlinks(t *testing.T)
{
dir := t.TempDir()
store, err := NewJobStore(dir)
if err != nil {
t.Fatal(err)
}
if err := store.Save(&JobRecord{Name: "../../escape"}); err == nil {
t.Fatal("Save() accepted path traversal")
}
if _, err := store.Load("../../escape"); err == nil {
t.Fatal("Load() accepted path traversal")
}
outside := filepath.Join(t.TempDir(), "outside.json")
if err := os.WriteFile(outside, []byte(`{"name":"outside"}`), 0600); err != nil {
t.Fatal(err)
}
if err := os.Symlink(outside, filepath.Join(dir, "linked.json")); err != nil {
t.Fatal(err)
}
if _, err := store.Load("linked"); err == nil {
t.Fatal("Load() followed a symbolic link")
}
list, err := store.List()
if err != nil {
t.Fatal(err)
}
if len(list) != 0 {
t.Fatalf("List() returned symlink record: %v", list)
}
}
TestJobStoreUsesPrivateModeAndSizeLimit
Parameters
func TestJobStoreUsesPrivateModeAndSizeLimit(t *testing.T)
{
dir := t.TempDir()
store, err := NewJobStore(dir)
if err != nil {
t.Fatal(err)
}
if err := store.Save(&JobRecord{Name: "private"}); err != nil {
t.Fatal(err)
}
info, err := os.Stat(filepath.Join(dir, "private.json"))
if err != nil {
t.Fatal(err)
}
if info.Mode().Perm() != 0600 {
t.Fatalf("file mode = %o, want 600", info.Mode().Perm())
}
if err := store.Save(&JobRecord{
Name: "large",
Cron: string(make([]byte, maxJobRecordSize+1)),
}); err == nil {
t.Fatal("Save() accepted an oversized record")
}
}
TestScheduler_WithStore
Parameters
func TestScheduler_WithStore(t *testing.T)
{
dir := t.TempDir()
s := New(WithStore(dir))
var count int64
if err := s.Register(Job{
Name: "tick",
Cron: "* * * * *",
Handler: func(ctx context.Context) error {
atomic.AddInt64(&count, 1)
return nil
},
}); err != nil {
t.Fatal(err)
}
s.runDue(context.Background(), time.Now())
time.Sleep(100 * time.Millisecond)
if atomic.LoadInt64(&count) == 0 {
t.Error("expected job to run at least once")
}
}
TestSchedulerRegisterSurfacesCorruptState
Parameters
func TestSchedulerRegisterSurfacesCorruptState(t *testing.T)
{
dir := t.TempDir()
if err := os.WriteFile(filepath.Join(dir, "job.json"), []byte("{"), 0600); err != nil {
t.Fatal(err)
}
scheduler := New(WithStore(dir))
err := scheduler.Register(Job{
Name: "job",
Cron: "* * * * *",
Handler: func(context.Context) error { return nil },
})
if err == nil {
t.Fatal("Register() ignored corrupt persisted state")
}
}
TestSchedulerPersistsHandlerErrorStatus
Parameters
func TestSchedulerPersistsHandlerErrorStatus(t *testing.T)
{
dir := t.TempDir()
scheduler := New(WithStore(dir))
if err := scheduler.Register(Job{
Name: "failing",
Cron: "* * * * *",
Handler: func(context.Context) error { return errors.New("failed") },
}); err != nil {
t.Fatal(err)
}
scheduler.runDue(context.Background(), time.Now())
deadline := time.Now().Add(time.Second)
for {
record, err := scheduler.store.Load("failing")
if err == nil {
if record.LastStatus != "error" {
t.Fatalf("LastStatus = %q, want error", record.LastStatus)
}
break
}
if time.Now().After(deadline) {
t.Fatalf("job state was not persisted: %v", err)
}
time.Sleep(time.Millisecond)
}
}
TestSchedulerLogsPersistenceFailure
Parameters
func TestSchedulerLogsPersistenceFailure(t *testing.T)
{
dir := t.TempDir()
messages := make(chan string, 10)
scheduler := New(WithStore(dir), WithLogger(func(message string) {
messages <- message
}))
if err := scheduler.Register(Job{
Name: "job",
Cron: "* * * * *",
Handler: func(context.Context) error { return nil },
}); err != nil {
t.Fatal(err)
}
if err := os.RemoveAll(dir); err != nil {
t.Fatal(err)
}
scheduler.runDue(context.Background(), time.Now())
timeout := time.After(time.Second)
for {
select {
case message := <-messages:
if strings.Contains(message, "persist job job state") {
if scheduler.Err() == nil {
t.Fatal("persistence failure was not exposed by Err()")
}
return
}
case <-timeout:
t.Fatal("persistence failure was not logged")
}
}
}