Documentation
¶
Overview ¶
Package scheduler provides a periodic task scheduler with minimal external dependencies (only github.com/ieshan/idx for ID generation). It supports cron expressions, fixed intervals, and one-shot schedules, with configurable concurrency, jittered exponential backoff, and graceful shutdown.
Key types:
- Scheduler — orchestrates job dispatch
- Job — a unit of scheduled work with its state and config
- JobExecutor — the interface callers implement to run jobs
- JobStore — persistence interface (in-memory or custom)
See https://github.com/ieshan/scheduler for guides.
Index ¶
- Variables
- func RetryBackoff(base time.Duration, attempt int) time.Duration
- type Config
- type DeliveryService
- type FileJobStore
- func (s *FileJobStore) Delete(_ context.Context, id idx.ID) error
- func (s *FileJobStore) Get(_ context.Context, id idx.ID) (*Job, error)
- func (s *FileJobStore) List(_ context.Context) ([]Job, error)
- func (s *FileJobStore) Save(_ context.Context, job *Job) error
- func (s *FileJobStore) UpdateState(_ context.Context, id idx.ID, state JobState) error
- type InMemoryJobStore
- func (s *InMemoryJobStore) Delete(_ context.Context, id idx.ID) error
- func (s *InMemoryJobStore) Get(_ context.Context, id idx.ID) (*Job, error)
- func (s *InMemoryJobStore) List(_ context.Context) ([]Job, error)
- func (s *InMemoryJobStore) Save(_ context.Context, job *Job) error
- func (s *InMemoryJobStore) UpdateState(_ context.Context, id idx.ID, state JobState) error
- type Job
- type JobConfig
- type JobExecutor
- type JobResult
- type JobState
- type JobStatus
- type JobStore
- type MessageSender
- type Option
- type RouterDelivery
- type Schedule
- type Scheduler
Examples ¶
Constants ¶
This section is empty.
Variables ¶
var ( // ErrJobNotFound indicates the requested job ID does not exist in the store. ErrJobNotFound = errors.New("job not found") )
Sentinel errors returned by JobStore implementations.
Functions ¶
func RetryBackoff ¶
RetryBackoff calculates a jittered exponential backoff duration. This is commonly used by JobExecutor implementations to determine how long to wait before retrying a failed job.
Formula: base * 2^attempt * (0.5 + rand(0, 1))
For example, with base=1s:
- attempt 0: 0.5s – 1.5s
- attempt 1: 1s – 3s
- attempt 2: 2s – 6s
The jitter prevents thundering herd when multiple jobs retry simultaneously.
Example ¶
ExampleRetryBackoff demonstrates jittered exponential backoff calculation.
base := time.Second
for i := range 3 {
d := RetryBackoff(base, i)
fmt.Printf("attempt %d: %v <= delay < %v\n", i, base*time.Duration(1<<uint(i))/2, base*time.Duration(1<<uint(i))*3/2)
_ = d // actual value varies due to jitter
}
Output: attempt 0: 500ms <= delay < 1.5s attempt 1: 1s <= delay < 3s attempt 2: 2s <= delay < 6s
Types ¶
type Config ¶
type Config struct {
// Store provides job persistence and query. Required.
Store JobStore
// Executors maps executor type names (e.g., "agent", "shell") to their [JobExecutor] implementations.
// Required. The scheduler uses this map to dispatch jobs to their executors.
Executors map[string]JobExecutor
// MaxConcurrent limits how many jobs execute simultaneously.
// Default: 5. Must be positive.
MaxConcurrent int
// PollInterval is the time between store polls.
// Default: 30s. Must be positive.
PollInterval time.Duration
// Logger receives scheduler operational messages (job dispatch, errors, state updates).
// Default: slog.Default().
Logger *slog.Logger
}
Config configures the Scheduler.
type DeliveryService ¶
type DeliveryService interface {
Deliver(ctx context.Context, result *JobResult) error
Close() error
}
DeliveryService routes completed job output to the originating channel.
type FileJobStore ¶ added in v1.0.1
type FileJobStore struct {
// contains filtered or unexported fields
}
FileJobStore is a file-based JobStore implementation that persists jobs to a JSON file. It provides thread-safe operations and atomic writes (write to temp file, then rename).
func NewFileJobStore ¶ added in v1.0.1
func NewFileJobStore(path string) (*FileJobStore, error)
NewFileJobStore creates a new FileJobStore that persists to the given file path. If the file exists, it will be loaded automatically.
func (*FileJobStore) Delete ¶ added in v1.0.1
Delete removes a job from the store. Returns ErrJobNotFound if the job does not exist.
func (*FileJobStore) Get ¶ added in v1.0.1
Get returns a single job by ID, or ErrJobNotFound if not found.
func (*FileJobStore) List ¶ added in v1.0.1
func (s *FileJobStore) List(_ context.Context) ([]Job, error)
List returns all jobs in the store.
func (*FileJobStore) Save ¶ added in v1.0.1
func (s *FileJobStore) Save(_ context.Context, job *Job) error
Save creates or overwrites a job in the store.
func (*FileJobStore) UpdateState ¶ added in v1.0.1
UpdateState updates only the JobState fields of a job. Returns ErrJobNotFound if the job does not exist.
type InMemoryJobStore ¶
type InMemoryJobStore struct {
// contains filtered or unexported fields
}
InMemoryJobStore is a thread-safe in-memory JobStore implementation for testing. It stores jobs in a map and provides full JobStore semantics.
This is useful for unit tests and demos where persistence is not required. For production, implement JobStore with a database backend.
Example ¶
ExampleInMemoryJobStore demonstrates basic CRUD operations on the in-memory store.
store := NewInMemoryJobStore()
ctx := context.Background()
job := &Job{
ID: mustID("01HZY0CWD0A0VKBQHHP3MS4GC3"),
Name: "demo-job",
State: JobState{
NextRun: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC),
},
}
_ = store.Save(ctx, job)
got, err := store.Get(ctx, mustID("01HZY0CWD0A0VKBQHHP3MS4GC3"))
if err != nil {
panic(err)
}
fmt.Println(got.Name)
jobs, _ := store.List(ctx)
fmt.Println(len(jobs))
_ = store.Delete(ctx, mustID("01HZY0CWD0A0VKBQHHP3MS4GC3"))
_, err = store.Get(ctx, mustID("01HZY0CWD0A0VKBQHHP3MS4GC3"))
fmt.Println(errors.Is(err, ErrJobNotFound))
Output: demo-job 1 true
func NewInMemoryJobStore ¶
func NewInMemoryJobStore() *InMemoryJobStore
NewInMemoryJobStore creates a new empty in-memory job store.
func (*InMemoryJobStore) Delete ¶
Delete removes a job from the store. Returns ErrJobNotFound if the job does not exist.
func (*InMemoryJobStore) Get ¶
Get returns a single job by ID, or ErrJobNotFound if not found.
func (*InMemoryJobStore) List ¶
func (s *InMemoryJobStore) List(_ context.Context) ([]Job, error)
List returns all jobs in the store.
func (*InMemoryJobStore) Save ¶
func (s *InMemoryJobStore) Save(_ context.Context, job *Job) error
Save creates or overwrites a job in the store.
func (*InMemoryJobStore) UpdateState ¶
UpdateState updates only the JobState fields of a job. Returns ErrJobNotFound if the job does not exist.
type Job ¶
type Job struct {
// ID is the unique identifier for this job.
ID idx.ID `json:"id"`
// Name is a human-readable name for the job.
Name string `json:"name"`
// Schedule is the parsed [Schedule] interface (Cron, Every, or At).
// Not serialized directly; the schedule type and expression are stored in ScheduleType.
Schedule Schedule `json:"-"` // serialized separately
// ScheduleType is the schedule type name ("cron", "every", "at").
// The scheduler uses this to reconstruct the Schedule interface on load.
ScheduleType string `json:"schedule_type"`
// ScheduleExpression is the raw schedule expression used for persistence.
// For "cron": "0 9 * * *", for "every": "5m", for "at": RFC3339 timestamp.
// This is used to reconstruct the Schedule interface when loading from storage.
ScheduleExpression string `json:"schedule_expression"`
// Payload is optional user-defined data passed to the executor.
Payload any `json:"payload"`
// Enabled indicates whether the job should be dispatched when it is due.
// The scheduler skips disabled jobs during tick.
Enabled bool `json:"enabled"`
// DeleteAfterRun indicates this is a one-shot job that should be disabled after successful execution.
DeleteAfterRun bool `json:"delete_after_run"`
// ExecutorType identifies which executor in the [Scheduler] should run this job
// (e.g., "agent", "shell", or a custom type).
ExecutorType string `json:"executor_type"` // "agent", "shell", etc.
// Config holds execution settings for this job (timeout, retries, backoff).
Config JobConfig `json:"config"`
// State holds runtime state: next scheduled run, last run time, status, output.
// Updated by the scheduler after each execution.
State JobState `json:"state"`
// ChannelKey identifies the delivery target (e.g. "tg:123", "tui:local").
// If set, the scheduler routes job results via [DeliveryService].
ChannelKey string `json:"channel_key,omitempty"`
// Prompt is the text sent to the agent executor.
// Used by [AgentJobExecutor] in the agent module.
Prompt string `json:"prompt,omitempty"`
// Script is the shell command for [ShellJobExecutor].
// Used by shell-based executors.
Script string `json:"script,omitempty"`
}
Job is a scheduled task with a schedule, executor type, payload, and runtime state. A job's state (NextRun, LastRun, RunCount) is updated by the scheduler after each execution.
type JobConfig ¶
type JobConfig struct {
// Timeout is the maximum duration a job is allowed to run.
// If zero, no timeout is enforced.
Timeout time.Duration `json:"timeout"`
// MaxRetries is the maximum number of times a retryable error triggers a retry.
// If zero, no retries are attempted (fail immediately on error).
MaxRetries int `json:"max_retries"`
// RetryBackoff is the base duration for jittered exponential backoff between retries.
// Actual backoff: [RetryBackoff] * 2^attempt * (0.5 + random).
RetryBackoff time.Duration `json:"retry_backoff"`
}
JobConfig holds per-job execution settings.
type JobExecutor ¶
type JobExecutor interface {
// Execute runs the job and returns the result or an error.
// The context is cancelled if the scheduler stops or the parent context expires.
Execute(ctx context.Context, job *Job) (*JobResult, error)
}
JobExecutor executes a job's payload. Consumers implement this interface to handle specific job types (e.g., agent jobs, shell scripts, webhooks).
Implementations must:
- Respect context cancellation
- Return a non-nil JobResult on success or an error on failure
- Handle their own timeout logic (job.Config.Timeout is informational)
- Distinguish retryable vs. deterministic errors if retry logic is needed
type JobResult ¶
type JobResult struct {
// Status is the outcome (success, failed, retrying, skipped).
Status JobStatus `json:"status"`
// Output is the execution output (stdout, response body, etc.).
Output string `json:"output"`
// Error is the error message if Status is failed or retrying.
// Empty if Status is success or skipped.
Error string `json:"error,omitempty"`
// Duration is the time taken to execute the job.
Duration time.Duration `json:"duration"`
// ChannelKey identifies the delivery target (e.g., "tg:123").
// Set by the executor if the job result should be routed back to a channel.
ChannelKey string `json:"channel_key,omitempty"` // delivery target
// Silent suppresses delivery when true.
// If set, the scheduler skips calling [DeliveryService.Deliver].
Silent bool `json:"silent,omitempty"`
}
JobResult is the outcome of a single job execution. Returned by JobExecutor.Execute.
type JobState ¶
type JobState struct {
// NextRun is the next scheduled time this job should execute.
// Calculated by the scheduler based on the job's Schedule.
NextRun time.Time `json:"next_run"`
// LastRun is the most recent time this job was executed.
LastRun time.Time `json:"last_run"`
// LastStatus is the outcome of the most recent execution (success, failed, retrying, skipped).
LastStatus JobStatus `json:"last_status"`
// LastOutput is the output or error message from the most recent execution.
LastOutput string `json:"last_output"`
// RunCount is the number of times this job has been executed.
RunCount int `json:"run_count"`
}
JobState holds runtime state for a job, updated after each execution.
type JobStatus ¶
type JobStatus string
JobStatus represents the outcome of a job execution.
const ( // StatusSuccess indicates the job completed successfully. StatusSuccess JobStatus = "success" // StatusFailed indicates the job encountered a non-retryable error. StatusFailed JobStatus = "failed" // StatusRetrying indicates the job will be retried (set by executor or scheduler retry logic). StatusRetrying JobStatus = "retrying" // StatusSkipped indicates the job was intentionally skipped (executor discretion). StatusSkipped JobStatus = "skipped" )
type JobStore ¶
type JobStore interface {
// List returns all jobs in the store.
List(ctx context.Context) ([]Job, error)
// Get returns a single job by ID, or an error if not found.
Get(ctx context.Context, id idx.ID) (*Job, error)
// Save creates or overwrites a job in the store.
// Returns an error if the job cannot be persisted.
Save(ctx context.Context, job *Job) error
// Delete removes a job from the store.
// Returns an error if the job is not found.
Delete(ctx context.Context, id idx.ID) error
// UpdateState updates only the [JobState] fields of a job
// (NextRun, LastRun, LastStatus, LastOutput, RunCount).
// The scheduler calls this after each execution.
// Returns an error if the job is not found.
UpdateState(ctx context.Context, id idx.ID, state JobState) error
}
JobStore persists Job definitions and their runtime JobState. Implementations provide in-memory or database-backed storage. The Scheduler calls these methods to load jobs, update state after execution, and persist changes.
type MessageSender ¶
MessageSender sends a text message to a target address via some transport.
type Option ¶
type Option func(*Scheduler)
Option configures a Scheduler after construction. Used with New.
func WithDelivery ¶
func WithDelivery(svc DeliveryService) Option
WithDelivery registers a DeliveryService that is called after each job completes. The service routes job results back to the originating channel (e.g., Telegram, TUI). If not set, results are discarded (but still updated in the store).
func WithNowFunc ¶
WithNowFunc overrides the clock used by the scheduler. By default, Scheduler uses time.Now. This is primarily useful for testing where a fake clock eliminates timing-dependent sleeps and makes tests deterministic.
Example (for tests):
sched := scheduler.New(cfg, scheduler.WithNowFunc(func() time.Time {
return time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC)
}))
type RouterDelivery ¶
type RouterDelivery struct {
// contains filtered or unexported fields
}
RouterDelivery routes based on the channel-key prefix ("tg:123" → prefix "tg", target "123").
func NewRouterDelivery ¶
func NewRouterDelivery(senders map[string]MessageSender, scanner func(string) string) *RouterDelivery
NewRouterDelivery creates a RouterDelivery. senders maps prefix → sender. scanner is an optional credential-redaction function applied before sending.
func (*RouterDelivery) Close ¶
func (r *RouterDelivery) Close() error
Close is a no-op for RouterDelivery.
func (*RouterDelivery) Deliver ¶
func (r *RouterDelivery) Deliver(ctx context.Context, result *JobResult) error
Deliver routes the result to the appropriate sender based on ChannelKey prefix. If result.Silent is true delivery is skipped. If result.Error is non-empty the error text is sent instead of the output.
type Schedule ¶
type Schedule interface {
// NextTick returns the next time this schedule fires after now.
// Returns zero time if the schedule will never fire again (e.g., a past [At] schedule).
NextTick(now time.Time) time.Time
// Type returns the schedule type name ("cron", "every", "at").
// Used during persistence and reconstruction.
Type() string
}
Schedule determines when a job should run. Implementations include Cron (5-field expression), Every (fixed interval), and At (one-shot).
func At ¶
At creates a one-shot Schedule that fires exactly once at time t. If t is in the past, Schedule.NextTick returns zero time and the job never fires. Typically used with Job.DeleteAfterRun to auto-disable after execution.
Example ¶
ExampleAt demonstrates a one-shot schedule.
target := time.Date(2026, 12, 25, 0, 0, 0, 0, time.UTC) s := At(target) // Before the target time — returns the target. now := time.Date(2026, 4, 15, 0, 0, 0, 0, time.UTC) fmt.Println(s.NextTick(now).Format(time.RFC3339)) // After the target time — returns zero time. later := time.Date(2027, 1, 1, 0, 0, 0, 0, time.UTC) fmt.Println(s.NextTick(later).IsZero())
Output: 2026-12-25T00:00:00Z true
func Cron ¶
Cron creates a cron Schedule from a 5-field expression: minute hour dom month dow. Returns an error if the expression has the wrong number of fields or values are out of range.
Examples:
- "0 9 * * *" — every day at 9:00 AM
- "0 9 * * 1-5" — weekdays at 9:00 AM
- "*/5 * * * *" — every 5 minutes
- "0 0 1 * *" — first day of each month at midnight
Example ¶
ExampleCron demonstrates parsing a cron expression and computing the next tick.
s, err := Cron("0 9 * * *")
if err != nil {
panic(err)
}
now := time.Date(2026, 4, 15, 8, 0, 0, 0, time.UTC)
next := s.NextTick(now)
fmt.Println(next.Format(time.RFC3339))
Output: 2026-04-15T09:00:00Z
func Every ¶
Every creates a fixed-interval Schedule that fires every duration. For example, Every(5*time.Minute) fires every 5 minutes.
Example ¶
ExampleEvery demonstrates a fixed-interval schedule.
s := Every(5 * time.Minute) now := time.Date(2026, 4, 15, 10, 0, 0, 0, time.UTC) next := s.NextTick(now) fmt.Println(next.Format(time.RFC3339))
Output: 2026-04-15T10:05:00Z
type Scheduler ¶
type Scheduler struct {
// contains filtered or unexported fields
}
Scheduler runs jobs on their schedules. It periodically polls the job store, identifies due jobs (where NextRun <= now), and dispatches them to registered JobExecutor implementations with bounded concurrency and graceful shutdown.
Create a scheduler with New, start it with Scheduler.Start, add jobs via the job store, and stop it with Scheduler.Stop.
Example:
store := scheduler.NewInMemoryJobStore()
cfg := scheduler.Config{
Store: store,
Executors: map[string]scheduler.JobExecutor{"shell": myShellExecutor},
MaxConcurrent: 5,
PollInterval: 30 * time.Second,
}
sched := scheduler.New(cfg)
go sched.Start(context.Background())
defer sched.Stop()
// Add a cron job to the store, then the scheduler will dispatch it...
Example (Cron) ¶
ExampleScheduler_cron demonstrates adding a cron job to a running Scheduler.
store := NewInMemoryJobStore()
sched := New(Config{
Store: store,
Executors: map[string]JobExecutor{"noop": &noopExecutor{}},
PollInterval: time.Second,
MaxConcurrent: 1,
})
ctx, cancel := context.WithCancel(context.Background())
go sched.Start(ctx)
defer sched.Stop()
defer cancel()
cronSched, err := Cron("0 9 * * *")
if err != nil {
panic(err)
}
job := &Job{
ID: mustID("01HZY0CWD0A0VKBQHHP3MS4GC1"),
Name: "daily-report",
ExecutorType: "noop",
Enabled: true,
Schedule: cronSched,
ScheduleType: "cron",
}
_ = store.Save(context.Background(), job)
saved, _ := store.Get(context.Background(), mustID("01HZY0CWD0A0VKBQHHP3MS4GC1"))
fmt.Println(saved.Name)
Output: daily-report
Example (Every) ¶
ExampleScheduler_every demonstrates adding a fixed-interval job.
store := NewInMemoryJobStore()
job := &Job{
ID: mustID("01HZY0CWD0A0VKBQHHP3MS4GC2"),
Name: "status-check",
ExecutorType: "noop",
Enabled: true,
Schedule: Every(5 * time.Minute),
ScheduleType: "every",
}
_ = store.Save(context.Background(), job)
saved, _ := store.Get(context.Background(), mustID("01HZY0CWD0A0VKBQHHP3MS4GC2"))
fmt.Println(saved.Name)
Output: status-check
func New ¶
New creates a new Scheduler from the given Config. The scheduler is not started; call Scheduler.Start to begin dispatching jobs. Options are applied after construction and may override config defaults.
Example ¶
ExampleNew demonstrates creating and starting a Scheduler.
store := NewInMemoryJobStore()
cfg := Config{
Store: store,
Executors: map[string]JobExecutor{"print": &printExecutor{}},
MaxConcurrent: 2,
PollInterval: 50 * time.Millisecond,
}
sched := New(cfg)
ctx, cancel := context.WithCancel(context.Background())
go sched.Start(ctx)
defer sched.Stop()
defer cancel()
job := &Job{
ID: mustID("01HZY0CWD0A0VKBQHHP3MS4GC0"),
Name: "hello-job",
ExecutorType: "print",
Enabled: true,
Schedule: Every(100 * time.Millisecond),
ScheduleType: "every",
State: JobState{NextRun: time.Now()},
}
_ = store.Save(context.Background(), job)
time.Sleep(250 * time.Millisecond)
func (*Scheduler) Start ¶
Start runs the scheduler's main loop. It blocks until Scheduler.Stop is called or ctx is cancelled. Start polls the job store at the configured Config.PollInterval, identifies due jobs (NextRun <= now), and dispatches them to their JobExecutor implementations.
Start is safe to call multiple times but will only start the scheduler on the first call. Typically called in a goroutine:
go sched.Start(context.Background())
func (*Scheduler) Stop ¶
func (s *Scheduler) Stop()
Stop gracefully stops the scheduler and waits for all in-flight jobs to complete. The context passed to Scheduler.Start is cancelled to signal executors to stop, and Stop blocks until all goroutines have exited.
After Stop returns, no further jobs will be dispatched. Stop is safe to call multiple times (idempotent). Stop is safe to call before Start() (no-op).
func (*Scheduler) Wake ¶
func (s *Scheduler) Wake()
Wake signals the scheduler to re-evaluate jobs immediately, without waiting for the next Config.PollInterval. This is useful when a job is added to the store and you want it to run ASAP. Wake is safe to call concurrently and does not block.