Documentation
¶
Overview ¶
Package cron provides durable in-process scheduling with human-readable declarations and explicit state stores.
Index ¶
- Constants
- Variables
- func OccurrenceID(jobName string, scheduledAt time.Time) string
- type CatchUpPolicy
- type ClaimDisposition
- type ClaimRequest
- type ClaimResult
- type Completion
- type IStore
- type Invocation
- type JobDefinition
- type JobFunc
- type JobOption
- type JobState
- type LeaseRenewal
- type MemoryStore
- func (store *MemoryStore) Claim(ctx context.Context, request ClaimRequest) (ClaimResult, error)
- func (store *MemoryStore) Complete(ctx context.Context, completion Completion) error
- func (store *MemoryStore) Expired(ctx context.Context, now time.Time, limit int) ([]OccurrenceRecord, error)
- func (store *MemoryStore) Reconcile(ctx context.Context, definitions []JobDefinition, now time.Time) ([]JobState, error)
- func (store *MemoryStore) Renew(ctx context.Context, renewal LeaseRenewal) error
- func (store *MemoryStore) Skip(ctx context.Context, request SkipRequest) error
- func (store *MemoryStore) Snapshot(ctx context.Context) (StoreSnapshot, error)
- type MeridiemTime
- type OccurrenceRecord
- type OccurrenceStatus
- type Option
- type OverlapPolicy
- type Rule
- type Schedule
- func (schedule Schedule) Anchor(anchor time.Time) Schedule
- func (schedule Schedule) Canonical() (string, error)
- func (schedule Schedule) Definition() (ScheduleDefinition, error)
- func (schedule Schedule) In(location *time.Location) Schedule
- func (schedule Schedule) IntervalAnchor() (time.Time, bool)
- func (schedule Schedule) Location() *time.Location
- func (schedule Schedule) Next(after time.Time) (time.Time, error)
- func (schedule Schedule) String() string
- func (schedule Schedule) Validate() error
- type ScheduleDefinition
- type ScheduleKind
- type Service
- type SkipRequest
- type Snapshot
- type StoreSnapshot
- type TimeOfDay
Constants ¶
const ( // SnapshotOccurrenceLimit bounds the recent occurrence history returned by // IStore.Snapshot and the read-only status route. SnapshotOccurrenceLimit = 1_000 )
const StatusPath = "/cron"
StatusPath is the sole read-only route mounted by Register.
Variables ¶
var ( // ErrInvalidConfiguration reports a constructor or persisted-definition // value that cannot form a safe scheduler. ErrInvalidConfiguration = errors.New("cron: invalid configuration") // ErrStoreRequired reports a Service constructed without an explicit IStore. ErrStoreRequired = errors.New("cron: store is required") // ErrNoJobs reports a Service constructed without any static jobs. ErrNoJobs = errors.New("cron: at least one job is required") // ErrJobNotFound reports a store operation for a job absent from the latest // static reconciliation. ErrJobNotFound = errors.New("cron: job not found") // ErrLeaseLost reports a stale or expired lease token. It is a fencing error: // the caller must not continue acting as the occurrence owner. ErrLeaseLost = errors.New("cron: occurrence lease lost") // ErrOccurrenceConflict reports reused occurrence identity with different // job or scheduled-time data. ErrOccurrenceConflict = errors.New("cron: occurrence identity conflict") // ErrOccurrenceRunning reports an attempt to skip an occurrence that is // already executing. ErrOccurrenceRunning = errors.New("cron: occurrence is running") // ErrAlreadyRunning reports concurrent Run calls on one Service. ErrAlreadyRunning = errors.New("cron: service is already running") )
Functions ¶
Types ¶
type CatchUpPolicy ¶
type CatchUpPolicy string
CatchUpPolicy controls which occurrences missed while the service was not running are considered when it starts again.
const ( // CatchUpNone advances past every startup occurrence without invoking it. // This is the conservative default and matches traditional cron behavior. CatchUpNone CatchUpPolicy = "none" // CatchUpLatest invokes only the latest startup occurrence. CatchUpLatest CatchUpPolicy = "latest" // CatchUpAll invokes every startup occurrence, up to the service catch-up // limit. Overlap policy is still enforced while those invocations run. CatchUpAll CatchUpPolicy = "all" )
type ClaimDisposition ¶
type ClaimDisposition string
ClaimDisposition explains the durable outcome of Claim.
const ( ClaimAcquired ClaimDisposition = "acquired" ClaimAlreadyTerminal ClaimDisposition = "already_terminal" ClaimLeaseHeld ClaimDisposition = "lease_held" ClaimSkippedOverlap ClaimDisposition = "skipped_overlap" )
type ClaimRequest ¶
type ClaimRequest struct {
OccurrenceID string
JobName string
ScheduledAt time.Time
NextRunAt time.Time
LeaseOwner string
LeaseToken string
ClaimedAt time.Time
LeaseUntil time.Time
}
ClaimRequest atomically advances the job cursor and attempts to acquire the occurrence's fenced lease.
type ClaimResult ¶
type ClaimResult struct {
Disposition ClaimDisposition
Occurrence OccurrenceRecord
}
ClaimResult is idempotent for a deterministic occurrence ID.
type Completion ¶
type Completion struct {
OccurrenceID string
LeaseToken string
Status OccurrenceStatus
FinishedAt time.Time
Error string
}
Completion terminally records one acquired invocation.
type IStore ¶
type IStore interface {
Reconcile(ctx context.Context, definitions []JobDefinition, now time.Time) ([]JobState, error)
Claim(ctx context.Context, request ClaimRequest) (ClaimResult, error)
Renew(ctx context.Context, renewal LeaseRenewal) error
Complete(ctx context.Context, completion Completion) error
Skip(ctx context.Context, request SkipRequest) error
Expired(ctx context.Context, now time.Time, limit int) ([]OccurrenceRecord, error)
Snapshot(ctx context.Context) (StoreSnapshot, error)
}
IStore is the durable scheduler boundary. Implementations must make Claim and Skip atomic with the associated NextRunAt advance. Complete and Renew must fence on LeaseToken. Reconcile must replace the active static definition set, preserve an already established interval anchor, and leave abandoned runs discoverable through Expired without rewinding the normal job cursor.
type Invocation ¶
type Invocation struct {
ID string `json:"id"`
JobName string `json:"job_name"`
ScheduledAt time.Time `json:"scheduled_at"`
StartedAt time.Time `json:"started_at"`
Attempt uint32 `json:"attempt"`
}
Invocation identifies the logical occurrence passed to a job handler.
type JobDefinition ¶
type JobDefinition struct {
Name string `json:"name"`
Schedule ScheduleDefinition `json:"schedule"`
CatchUp CatchUpPolicy `json:"catch_up"`
Overlap OverlapPolicy `json:"overlap"`
}
JobDefinition is the static, persistence-neutral declaration reconciled at startup. Adapters map it to their own SQLC or wire types at the boundary.
type JobFunc ¶
type JobFunc func(ctx context.Context, invocation Invocation) error
JobFunc executes one scheduled invocation. It should stop promptly when ctx is cancelled. Returning an error records a failed occurrence; it does not stop the scheduler.
type JobOption ¶
type JobOption func(config *jobConfig) error
JobOption configures one job registered with WithJob.
func WithCatchUp ¶
func WithCatchUp(policy CatchUpPolicy) JobOption
WithCatchUp sets a job's missed-occurrence policy.
func WithOverlap ¶
func WithOverlap(policy OverlapPolicy) JobOption
WithOverlap sets a job's concurrent-occurrence policy.
type JobState ¶
type JobState struct {
Definition JobDefinition `json:"definition"`
NextRunAt time.Time `json:"next_run_at"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
JobState is the durable scheduling cursor for one active definition.
type LeaseRenewal ¶
type LeaseRenewal struct {
OccurrenceID string
LeaseToken string
RenewedAt time.Time
LeaseUntil time.Time
}
LeaseRenewal extends a live lease while preserving its fencing token.
type MemoryStore ¶
type MemoryStore struct {
// contains filtered or unexported fields
}
MemoryStore is an explicit process-local IStore for tests and disposable services. It implements the same lease and reconciliation semantics as a durable adapter, but intentionally does not survive process restart.
func NewMemoryStore ¶
func NewMemoryStore() *MemoryStore
NewMemoryStore returns an empty process-local store.
func (*MemoryStore) Claim ¶
func (store *MemoryStore) Claim(ctx context.Context, request ClaimRequest) (ClaimResult, error)
Claim implements atomically fenced occurrence acquisition for MemoryStore.
func (*MemoryStore) Complete ¶
func (store *MemoryStore) Complete(ctx context.Context, completion Completion) error
Complete implements fenced terminal recording for MemoryStore.
func (*MemoryStore) Expired ¶
func (store *MemoryStore) Expired(ctx context.Context, now time.Time, limit int) ([]OccurrenceRecord, error)
Expired returns a bounded stable-order view of abandoned running occurrences for active jobs. It does not mutate leases or scheduling cursors; callers recover records through Claim's normal fencing path.
func (*MemoryStore) Reconcile ¶
func (store *MemoryStore) Reconcile( ctx context.Context, definitions []JobDefinition, now time.Time, ) ([]JobState, error)
Reconcile implements static startup reconciliation for MemoryStore.
func (*MemoryStore) Renew ¶
func (store *MemoryStore) Renew(ctx context.Context, renewal LeaseRenewal) error
Renew implements fenced lease renewal for MemoryStore.
func (*MemoryStore) Skip ¶
func (store *MemoryStore) Skip(ctx context.Context, request SkipRequest) error
Skip implements idempotent terminal skip recording for MemoryStore.
func (*MemoryStore) Snapshot ¶
func (store *MemoryStore) Snapshot(ctx context.Context) (StoreSnapshot, error)
Snapshot returns a defensive, stable-order copy of MemoryStore state.
type MeridiemTime ¶
type MeridiemTime struct {
// contains filtered or unexported fields
}
MeridiemTime is deliberately distinct from TimeOfDay. It can only become a schedule time after AM or PM is selected, preventing Daily(At(3)).
func At ¶
func At(hour int, minute ...int) MeridiemTime
At starts a 12-hour clock declaration. It accepts an optional minute.
func (MeridiemTime) AM ¶
func (value MeridiemTime) AM() TimeOfDay
AM completes a 12-hour declaration.
func (MeridiemTime) PM ¶
func (value MeridiemTime) PM() TimeOfDay
PM completes a 12-hour declaration.
type OccurrenceRecord ¶
type OccurrenceRecord struct {
ID string `json:"id"`
JobName string `json:"job_name"`
ScheduledAt time.Time `json:"scheduled_at"`
Status OccurrenceStatus `json:"status"`
Attempt uint32 `json:"attempt"`
StartedAt time.Time `json:"started_at,omitempty"`
FinishedAt time.Time `json:"finished_at,omitempty"`
LeaseOwner string `json:"lease_owner,omitempty"`
LeaseToken string `json:"-"`
LeaseUntil time.Time `json:"lease_until,omitempty"`
Error string `json:"error,omitempty"`
SkipReason string `json:"skip_reason,omitempty"`
LastModified time.Time `json:"last_modified"`
}
OccurrenceRecord is the durable execution record for one scheduled instant. LeaseToken is a fencing secret used only by IStore implementations and is deliberately omitted from JSON snapshots.
type OccurrenceStatus ¶
type OccurrenceStatus string
OccurrenceStatus is the durable terminal or running state of an occurrence.
const ( OccurrenceRunning OccurrenceStatus = "running" OccurrenceSucceeded OccurrenceStatus = "succeeded" OccurrenceFailed OccurrenceStatus = "failed" OccurrenceCanceled OccurrenceStatus = "canceled" OccurrenceSkipped OccurrenceStatus = "skipped" )
type Option ¶
type Option func(config *serviceConfig) error
Option configures a Service. New validates the complete option set before reconciling or starting any work.
func WithCatchUpLimit ¶
WithCatchUpLimit bounds the number of due occurrences one job may process in a single scheduling cycle.
func WithLeaseDuration ¶
WithLeaseDuration sets how long an invocation lease remains valid without a renewal. Active jobs renew their lease three times per duration.
func WithLeaseOwner ¶
WithLeaseOwner sets the process identity recorded on leases. Most callers should use the random identity generated by New; this option is useful when an operator already has a stable, unique replica identity.
type OverlapPolicy ¶
type OverlapPolicy string
OverlapPolicy controls whether two occurrences of one job may execute at the same time. It is enforced by IStore, so it also covers multiple processes sharing a durable store.
const ( // OverlapSkip records a skipped occurrence when another occurrence of the // same job owns a live lease. This is the default. OverlapSkip OverlapPolicy = "skip" // OverlapAllow permits concurrent occurrences of the same job. OverlapAllow OverlapPolicy = "allow" )
type Rule ¶
type Rule struct {
// contains filtered or unexported fields
}
Rule is an opaque schedule declaration constructed by the fluent helpers. It intentionally has no exported fields so a caller cannot construct an unvalidated wire-shaped schedule.
func Every ¶
Every returns an interval rule. Its cadence is anchored when Anchor is supplied by the durable runtime; absent an explicit anchor, Next uses the Unix epoch as a stable default rather than process start time.
func LastDayOfMonth ¶
LastDayOfMonth returns a rule that fires on each month's final local day.
func Monthly ¶
Monthly returns a rule that fires on day (1 through 31). Months without that day are skipped.
type Schedule ¶
type Schedule struct {
// contains filtered or unexported fields
}
Schedule is an immutable, validated-on-use scheduling definition. A zero Schedule is invalid. Its location defaults to UTC unless In is called.
Typed-rule DST semantics are deliberately instant-based: Next walks real UTC minutes and matches their local civil representation. Spring-forward civil times do not run because they do not exist; a matching repeated fall-back time runs once for each real instant. Raw rules use robfig/cron's semantics.
func ScheduleFromDefinition ¶
func ScheduleFromDefinition(definition ScheduleDefinition) (Schedule, error)
ScheduleFromDefinition reconstructs and validates a Schedule from its neutral domain projection. It rejects a mismatched Canonical value so an adapter cannot silently reinterpret persisted schedule fields.
func Spec ¶
Spec wraps a rule in a Schedule. It does not panic: invalid declarations are retained and reported by Validate, Canonical, or Next.
func (Schedule) Anchor ¶
Anchor returns a copy with a durable interval anchor. It applies only to Every rules and lets a runtime retain cadence across restarts. The anchor is normalized to PostgreSQL's microsecond timestamp precision.
func (Schedule) Canonical ¶
Canonical returns a normalized five-field cron expression, or @every for an interval. The separate location and anchor are intentionally not encoded in that expression.
func (Schedule) Definition ¶
func (schedule Schedule) Definition() (ScheduleDefinition, error)
Definition returns a structured, normalized, persistence-ready projection.
func (Schedule) In ¶
In returns a copy scheduled in location. A nil location is reported by Validate instead of panicking.
func (Schedule) IntervalAnchor ¶
IntervalAnchor reports the explicit durable anchor, if any.
func (Schedule) Location ¶
Location returns the schedule location after validation. UTC is the default.
func (Schedule) Next ¶
Next returns the first matching real instant strictly after after. Raw rules delegate five-field parsing and occurrence evaluation to robfig/cron. Typed calendar rules are bounded to five years so malformed state cannot spin.
type ScheduleDefinition ¶
type ScheduleDefinition struct {
Kind ScheduleKind `json:"kind"`
Timezone string `json:"timezone"`
Canonical string `json:"canonical"`
Hour int `json:"hour,omitempty"`
Minute int `json:"minute,omitempty"`
Weekday time.Weekday `json:"weekday,omitempty"`
MonthDay int `json:"month_day,omitempty"`
Interval time.Duration `json:"interval,omitempty"`
Anchor time.Time `json:"anchor,omitempty"`
HasAnchor bool `json:"has_anchor,omitempty"`
}
ScheduleDefinition is the neutral persistence and boundary projection of a Schedule. It deliberately contains ordinary typed Go values: adapters map this into relational SQLC parameters or API messages at their own boundary. Canonical is a normalized expression, while the typed fields preserve the pleasant DSL shape for schedule kinds that have one.
type ScheduleKind ¶
type ScheduleKind string
ScheduleKind identifies the structured scheduling form in a Definition. It is a domain value, deliberately independent of protobuf and SQLC types.
const ( ScheduleKindDaily ScheduleKind = "daily" ScheduleKindWeekly ScheduleKind = "weekly" ScheduleKindMonthly ScheduleKind = "monthly" ScheduleKindLastDayOfMonth ScheduleKind = "last_day_of_month" ScheduleKindEvery ScheduleKind = "every" ScheduleKindRaw ScheduleKind = "raw" )
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
Service reconciles static jobs, schedules deterministic occurrences, and executes handlers under durable leases.
func New ¶
New constructs a stopped Service. An IStore and at least one WithJob are required; construction performs no persistence and starts no goroutines.
func (*Service) Register ¶
Register mounts the scheduler's read-only JSON status route. It never mounts mutation or run-now endpoints; callers choose the Gin group and therefore the authentication boundary.
type SkipRequest ¶
type SkipRequest struct {
OccurrenceID string
JobName string
ScheduledAt time.Time
NextRunAt time.Time
SkippedAt time.Time
Reason string
}
SkipRequest terminally records a deliberately uninvoked occurrence and atomically advances its job cursor.
type Snapshot ¶
type Snapshot struct {
GeneratedAt time.Time `json:"generated_at"`
Running bool `json:"running"`
Jobs []JobState `json:"jobs"`
Occurrences []OccurrenceRecord `json:"occurrences"`
}
Snapshot is the read-only service and durable-store projection returned by Snapshot and the Gin status route.
type StoreSnapshot ¶
type StoreSnapshot struct {
Jobs []JobState `json:"jobs"`
Occurrences []OccurrenceRecord `json:"occurrences"`
}
StoreSnapshot is a point-in-time copy of active jobs and at most the most recent SnapshotOccurrenceLimit durable occurrences.
Directories
¶
| Path | Synopsis |
|---|---|
|
Package contract maps the cron domain model to validated Liquid Proto messages at HTTP and messaging boundaries.
|
Package contract maps the cron domain model to validated Liquid Proto messages at HTTP and messaging boundaries. |
|
Package postgres provides the SQLC-backed durable cron Store.
|
Package postgres provides the SQLC-backed durable cron Store. |