cron

package
v0.2.4 Latest Latest
Warning

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

Go to latest
Published: Oct 3, 2026 License: Apache-2.0 Imports: 16 Imported by: 0

README

Cron: the trigger scheduler

services/cron is CSF's cron service: a durable in-process scheduler that mounts into the host runtime, fires each declared trigger on its schedule and records every occurrence through csfpg. The schedule grammar, the human-readable declarations and their canonical five-field form, is the pure library pkg/cron; this package is the behavior.

A trigger is a name, a schedule, a catch-up policy, an overlap policy and the operation it invokes. Each firing is an occurrence, identified by the trigger and its scheduled instant and run under a fenced lease. Execution is at least once: a lease can expire after the operation produced an external effect but before its end was recorded, so an operation uses the occurrence ID as its idempotency key.

Mounting

The binary opens the pool, applies CSF's schema and grants the service its store, its clock and its triggers; the host runtime starts it, stops it and joins it with every other service.

import (
	"context"
	"time"

	"github.com/candacelabs/csf/ipc/clock"
	"github.com/candacelabs/csf/ipc/db/csfpg"
	"github.com/candacelabs/csf/pkg/cron"
	"github.com/candacelabs/csf/runtime"
	cronservice "github.com/candacelabs/csf/services/cron"
)

func run(ctx context.Context, settings csfpg.Settings) error {
	// The binary owns the pool and CSF's schema; the service borrows the pool.
	pool, err := csfpg.OpenPool(ctx, settings)
	if err != nil {
		return err
	}
	defer pool.Close()
	handle := pool.OpenSQL()
	err = csfpg.ApplySchema(ctx, handle)
	_ = handle.Close()
	if err != nil {
		return err
	}
	store, err := cronservice.NewStore(pool)
	if err != nil {
		return err
	}
	scheduler, err := cronservice.NewScheduler(
		cronservice.WithStore(store),
		cronservice.WithClock(clock.NewSystemClock()),
		cronservice.WithTrigger("daily-rollup", cron.Spec(cron.Daily(cron.At(3).AM())), buildDailyRollup,
			cronservice.WithCatchUp(cron.CatchUpLatest)),
		cronservice.WithTrigger("cache-refresh", cron.Spec(cron.Every(15*time.Minute)), refreshCache,
			cronservice.WithOverlap(cron.OverlapSkip)),
	)
	if err != nil {
		return err
	}
	host, err := runtime.NewHostRuntime(runtime.WithHostName("rollups"))
	if err != nil {
		return err
	}
	if err := host.Mount("cron", scheduler); err != nil {
		return err
	}
	// Run returns once the scheduler has stopped and the occurrence in
	// flight, canceled and recorded as such, has been joined.
	return host.Run(ctx)
}

func buildDailyRollup(ctx context.Context, occurrence cronservice.Occurrence) error {
	// occurrence.ID is the idempotency key for this scheduled instant.
	return nil
}

Start reconciles the declared triggers with the store, so a store that cannot be read fails the mount, then starts the scheduling goroutine on the scope the runtime hands it. Every occurrence runs on a goroutine of that same scope: canceling the scope cancels the operation in flight, and joining it waits until that occurrence has recorded its end. A store that refuses a record afterwards fails the scope, and so the runtime, with the store's error as the cause; an operation's own error or panic is recorded as a failed occurrence and never stops the service.

Policies and their defaults

Both defaults are the conservative choice, and both are per trigger:

Trigger option Values Default Effect
WithCatchUp CatchUpNone · CatchUpLatest · CatchUpAll CatchUpNone What to do with occurrences missed while the process was down: skip past all of them (traditional cron), run only the most recent, or run every one up to the catch-up limit.
WithOverlap OverlapSkip · OverlapAllow OverlapSkip Whether a second occurrence may run while another holds a live lease. Enforced by the store, so it holds across processes sharing one database, not just within one.

Scheduler-wide options: WithStore (required), WithClock (the host's clock by default), WithLeaseDuration (30s; a running occurrence renews three times per duration), WithCatchUpLimit (1,000 due occurrences per trigger per cycle), and WithLeaseOwner for a binary that already has a stable replica identity; NewScheduler generates a random one otherwise.

NewScheduler rejects a duplicate trigger name, a name outside ^[a-z][a-z0-9._/-]*$, an invalid schedule, a nil operation, a missing store and an empty trigger set before anything is mounted.

The store

Cron state is two tables of CSF's one migration, ipc/db/csfpg/schema/001_init.sql: csf_cron_triggers, one row per declared trigger with its schedule as relational columns and its cursor, and csf_cron_occurrences, one row per scheduled instant with its status, attempt and lease. The queries are ipc/db/csfpg/cron.sql, generated by csfpg's sqlc set; NewStore runs them over the csfpg.IDB capability the binary grants. Every fenced write is one conditional statement, so the store takes no row lock or advisory lock and the same SQL runs on pgmem.

IStore is the scheduler's boundary. Claim acquires an occurrence's lease and advances the trigger's cursor in one transaction and is idempotent on the occurrence ID; Renew and Complete fence on the lease token; Skip records an occurrence a policy chose not to run; Expired lists abandoned leases for the scheduler to reclaim; Snapshot is the active triggers and the 1,000 most recent occurrences. The Liquid Proto messages under pkg/cron/v1 are portable boundary contracts, never a storage format.

Specs

Specs grant the service a clock.ManualClock and the store over pgmem, and move time themselves: crontest.OpenStore(GinkgoT()) applies CSF's real schema to a pgmem database and returns the store, closed when the spec ends. The suite proves, without a sleep, that a trigger fires at its instant, that an overlapping occurrence is skipped, that each catch-up policy does what it says after downtime, and that a host runtime's shutdown joins the occurrence in flight. pgmem keeps timestamps at second precision, so its specs use whole seconds.

The opt-in acceptance tier runs the store's SQL on a disposable PostgreSQL, opened only through csfpg.OpenPool:

CANDACE_CSF_TEST_DATABASE_URL='postgresql://csf:csf@localhost:5432/csf_test?sslmode=disable' \
  go test -race -tags acceptance ./services/cron/

Documentation

Overview

Package cron is the cron service: the durable in-process scheduler that mounts into a host runtime, fires each declared trigger on its schedule and records every occurrence through csfpg. A trigger is a named, human-readable schedule (candace/pkg/cron) with a catch-up and an overlap policy that invokes one operation of a mounted service; each firing is an occurrence, run under a fenced lease so that execution is at least once and survives a restart. The store, the clock and the triggers are granted through options; Start starts every goroutine through the scope the runtime hands it, and the scope's join is the service's cleanup.

Index

Constants

View Source
const (
	// DefaultLeaseDuration is how long an occurrence's lease stays valid
	// without a renewal; a running occurrence renews three times per
	// duration.
	DefaultLeaseDuration = 30 * time.Second
	// DefaultCatchUpLimit bounds the due occurrences one trigger processes
	// in one scheduling cycle.
	DefaultCatchUpLimit = 1_000
)
View Source
const (

	// SnapshotOccurrenceLimit bounds the recent occurrence history a
	// snapshot carries.
	SnapshotOccurrenceLimit = 1_000
)

Variables

View Source
var (
	// ErrInvalidConfiguration reports an option, a declaration or a stored
	// value that cannot form a safe scheduler.
	ErrInvalidConfiguration = errors.New("cron: invalid configuration")
	// ErrStoreRequired reports a scheduler built without [WithStore].
	ErrStoreRequired = errors.New("cron: a store is required")
	// ErrNoTriggers reports a scheduler built without a [WithTrigger].
	ErrNoTriggers = errors.New("cron: at least one trigger is required")
	// ErrAlreadyStarted reports a second Start of one scheduler.
	ErrAlreadyStarted = errors.New("cron: scheduler is already started")
)
View Source
var (
	// ErrDatabaseRequired reports a store built without the database
	// capability.
	ErrDatabaseRequired = errors.New("cron: a database is required")
	// ErrTriggerNotFound reports a store operation for a trigger absent from
	// the latest reconciliation.
	ErrTriggerNotFound = errors.New("cron: trigger not found")
	// ErrLeaseLost reports a stale or expired lease token. It is a fencing
	// error: the caller must stop acting as the occurrence's owner.
	ErrLeaseLost = errors.New("cron: occurrence lease lost")
	// ErrOccurrenceConflict reports an occurrence identity reused with
	// different trigger or scheduled-time data, or a cursor the request does
	// not follow.
	ErrOccurrenceConflict = errors.New("cron: occurrence identity conflict")
	// ErrOccurrenceRunning reports an attempt to skip an occurrence that is
	// executing under a live lease.
	ErrOccurrenceRunning = errors.New("cron: occurrence is running")
)

Functions

This section is empty.

Types

type ClaimDisposition

type ClaimDisposition string

ClaimDisposition is the durable outcome of a Claim.

const (
	ClaimAcquired        ClaimDisposition = "acquired"
	ClaimAlreadyTerminal ClaimDisposition = "already_terminal"
	ClaimLeaseHeld       ClaimDisposition = "lease_held"
	ClaimSkippedOverlap  ClaimDisposition = "skipped_overlap"
)

type ClaimRequest

type ClaimRequest struct {
	OccurrenceID string
	TriggerName  string
	ScheduledAt  time.Time
	NextRunAt    time.Time
	LeaseOwner   string
	LeaseToken   string
	ClaimedAt    time.Time
	LeaseUntil   time.Time
}

ClaimRequest atomically advances the trigger's cursor and acquires the occurrence's fenced lease.

type ClaimResult

type ClaimResult struct {
	Disposition ClaimDisposition
	Occurrence  grammar.OccurrenceRecord
}

ClaimResult is idempotent for a deterministic occurrence ID.

type Completion

type Completion struct {
	OccurrenceID string
	LeaseToken   string
	Status       grammar.OccurrenceStatus
	FinishedAt   time.Time
	Error        string
}

Completion records the end of one acquired occurrence.

type IStore

type IStore interface {
	Reconcile(ctx context.Context, definitions []grammar.TriggerDefinition, now time.Time) ([]grammar.TriggerState, 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) ([]grammar.OccurrenceRecord, error)
	Snapshot(ctx context.Context) (grammar.StoreSnapshot, error)
}

IStore is the scheduler's durable boundary: the out-of-process store the triggers' cursors and occurrences live in. Claim and Skip advance the cursor atomically with the record they write; Complete and Renew fence on the lease token; Reconcile replaces the active declarations, keeps an established interval anchor, and leaves abandoned occurrences for Expired without rewinding a cursor.

type LeaseRenewal

type LeaseRenewal struct {
	OccurrenceID string
	LeaseToken   string
	RenewedAt    time.Time
	LeaseUntil   time.Time
}

LeaseRenewal extends a live lease while keeping its fencing token.

type Occurrence

type Occurrence struct {
	ID          string    `json:"id"`
	TriggerName string    `json:"trigger_name"`
	ScheduledAt time.Time `json:"scheduled_at"`
	StartedAt   time.Time `json:"started_at"`
	Attempt     uint32    `json:"attempt"`
}

Occurrence identifies the firing an operation is running: its stable ID, the idempotency key for any external effect, and the attempt number.

type Operation

type Operation func(ctx context.Context, occurrence Occurrence) error

Operation is what a trigger invokes: one operation of a mounted service. It stops promptly when ctx is canceled. An error records a failed occurrence and never stops the scheduler.

type Option

type Option func(configuration *configuration) error

Option configures a Scheduler; NewScheduler validates the whole set before building anything.

func WithCatchUpLimit

func WithCatchUpLimit(limit int) Option

WithCatchUpLimit bounds the due occurrences one trigger processes in one scheduling cycle.

func WithClock

func WithClock(source clock.IClock) Option

WithClock grants the clock the scheduler reads and waits on. The host's clock is the default; a spec grants a controllable one.

func WithLeaseDuration

func WithLeaseDuration(duration time.Duration) Option

WithLeaseDuration sets how long an occurrence's lease stays valid without a renewal.

func WithLeaseOwner

func WithLeaseOwner(owner string) Option

WithLeaseOwner names this process on the leases it holds. NewScheduler generates a random identity otherwise; a binary with a stable replica identity passes it here.

func WithStore

func WithStore(store IStore) Option

WithStore grants the store occurrences are recorded through: NewStore over the csfpg capability in a binary, the same store over pgmem in a spec. Required.

func WithTrigger

func WithTrigger(name string, schedule grammar.Schedule, operation Operation, options ...TriggerOption) Option

WithTrigger declares one trigger: a name, its human-readable schedule and the operation each occurrence invokes, with WithCatchUp and WithOverlap choosing its policies. At least one is required, and names are unique.

type Scheduler

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

Scheduler is the cron service. It reconciles its declared triggers with the store when started, fires each due occurrence under a lease on its scope, and records every outcome; a failure to read or write the store safely fails the scope, and so the runtime.

func NewScheduler

func NewScheduler(options ...Option) (*Scheduler, error)

NewScheduler validates the whole option set and returns a stopped scheduler: it persists nothing and starts no goroutine until mounted.

func (*Scheduler) Start

func (scheduler *Scheduler) Start(scope *runtime.Scope) error

Start reconciles the declared triggers with the store, so a store that cannot be read fails the mount, then starts the scheduling goroutine on scope. Every occurrence runs on a goroutine of the same scope: canceling the scope cancels the occurrences in flight, and joining it waits for each to record its end.

type SkipRequest

type SkipRequest struct {
	OccurrenceID string
	TriggerName  string
	ScheduledAt  time.Time
	NextRunAt    time.Time
	SkippedAt    time.Time
	Reason       string
}

SkipRequest records an occurrence deliberately not invoked and advances its trigger's cursor past it.

type Store

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

Store is the IStore over CSF's PostgreSQL schema, reached through the csfpg capability: the csf_cron_triggers and csf_cron_occurrences tables and the queries csfpg generates for them. Every fenced write is one conditional statement, so it takes no row lock and runs unchanged on pgmem. It borrows the capability and never closes it.

func NewStore

func NewStore(database csfpg.IDB) (*Store, error)

NewStore returns the store over a pool the binary opened through ipc/db/csfpg, or over pgmem's IDB in a spec.

func (*Store) Claim

func (store *Store) Claim(ctx context.Context, request ClaimRequest) (ClaimResult, error)

Claim acquires the occurrence's lease and advances the trigger's cursor in one transaction. It is idempotent on the occurrence ID: a terminal occurrence, a live lease and an overlap are reported, not repeated.

func (*Store) Complete

func (store *Store) Complete(ctx context.Context, completion Completion) error

Complete records the end of an acquired occurrence under its live lease.

func (*Store) Expired

func (store *Store) Expired(ctx context.Context, now time.Time, limit int) ([]grammar.OccurrenceRecord, error)

Expired lists the abandoned running occurrences of active triggers, oldest expiry first, without moving a cursor; Claim reclaims each.

func (*Store) Reconcile

func (store *Store) Reconcile(ctx context.Context, definitions []grammar.TriggerDefinition, now time.Time) ([]grammar.TriggerState, error)

Reconcile makes definitions the active set: a new trigger starts at its first occurrence after now, a changed schedule restarts from now, a changed policy keeps its cursor, and a trigger no longer declared is disabled with its history kept.

func (*Store) Renew

func (store *Store) Renew(ctx context.Context, renewal LeaseRenewal) error

Renew extends a live lease; a stale token or an expired lease is lost.

func (*Store) Skip

func (store *Store) Skip(ctx context.Context, request SkipRequest) error

Skip records an occurrence that is deliberately not invoked and advances the cursor past it. A replay of a terminal skip changes nothing, and an occurrence under a live lease is refused.

func (*Store) Snapshot

func (store *Store) Snapshot(ctx context.Context) (grammar.StoreSnapshot, error)

Snapshot is the active triggers and the newest SnapshotOccurrenceLimit occurrences, in a stable order.

type TriggerOption

type TriggerOption func(policies *triggerPolicies) error

TriggerOption configures one trigger declared with WithTrigger.

func WithCatchUp

func WithCatchUp(policy grammar.CatchUpPolicy) TriggerOption

WithCatchUp sets a trigger's missed-occurrence policy.

func WithOverlap

func WithOverlap(policy grammar.OverlapPolicy) TriggerOption

WithOverlap sets a trigger's concurrent-occurrence policy.

Directories

Path Synopsis
Package crontest opens the cron store on pgmem, CSF's in-process substitute for PostgreSQL, with CSF's real schema applied: the store a spec grants when it needs cron state without a database.
Package crontest opens the cron store on pgmem, CSF's in-process substitute for PostgreSQL, with CSF's real schema applied: the store a spec grants when it needs cron state without a database.
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.

Jump to

Keyboard shortcuts

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