scheduler

package
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: MIT Imports: 6 Imported by: 0

Documentation

Overview

Package scheduler delivers durable Timebox messages when they become due

Index

Constants

View Source
const (
	// DefaultRescanInterval controls discovery of changes from other processes
	DefaultRescanInterval = time.Second

	// DefaultRetryDelay controls local retry after scheduler errors
	DefaultRetryDelay = time.Second
)

Variables

View Source
var (
	// ErrStoreRequired indicates a Scheduler has no Store
	ErrStoreRequired = errors.New("scheduler store is required")

	// ErrProcessorRequired indicates a Scheduler has no Processor
	ErrProcessorRequired = errors.New("scheduler processor is required")

	// ErrRetry leaves a due message active without logging the attempt
	ErrRetry = errors.New("retry scheduled message")

	// ErrInvalidRescanInterval indicates a non-positive rescan interval
	ErrInvalidRescanInterval = errors.New(
		"scheduler rescan interval must be positive",
	)

	// ErrInvalidRetryDelay indicates a non-positive retry delay
	ErrInvalidRetryDelay = errors.New("scheduler retry delay must be positive")
)

Functions

This section is empty.

Types

type Clock

type Clock func() time.Time

Clock provides the current scheduler time

type Config

type Config struct {
	Store            *timebox.Store
	Processor        Processor
	Clock            Clock
	TimerConstructor TimerConstructor
	RescanInterval   time.Duration
	RetryDelay       time.Duration
}

Config configures a disposable schedule runner

type Processor

type Processor func(*timebox.Transaction, *timebox.Message) error

Processor changes state in response to one due message

type Scheduler

type Scheduler struct {
	// contains filtered or unexported fields
}

Scheduler maintains a disposable heap over durable schedule aggregates

func New

func New(cfg Config) (*Scheduler, error)

New constructs a disposable schedule runner

func (*Scheduler) Run

func (s *Scheduler) Run(ctx context.Context) error

Run delivers due messages until ctx ends

func (*Scheduler) Wake

func (s *Scheduler) Wake()

Wake requests prompt reconciliation with the durable schedule index

type Timer

type Timer interface {
	Channel() <-chan time.Time
	Reset(time.Duration) bool
	Stop() bool
}

Timer provides the resettable wakeup used by Scheduler

type TimerConstructor

type TimerConstructor func(time.Duration) Timer

TimerConstructor creates a Timer for a delay

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL