postgres

package
v0.1.3 Latest Latest
Warning

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

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

Documentation

Overview

Package postgres provides the SQLC-backed durable cron Store.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type AcquireExpiredRunParams

type AcquireExpiredRunParams struct {
	WorkerID     sql.NullString
	LeaseToken   sql.NullString
	LeaseUntil   sql.NullTime
	ClaimedAt    sql.NullTime
	OccurrenceID string
}

type AdvanceJobNextRunParams

type AdvanceJobNextRunParams struct {
	NextRunAt        sql.NullTime
	ScheduleCursorAt sql.NullTime
	UpdatedAt        time.Time
	JobName          string
}

type CandaceCronJob

type CandaceCronJob struct {
	JobName             string
	ScheduleKind        string
	LocalHour           sql.NullInt16
	LocalMinute         sql.NullInt16
	Weekday             sql.NullInt16
	MonthDay            sql.NullInt16
	IntervalNanoseconds sql.NullInt64
	RawExpression       sql.NullString
	Timezone            string
	IntervalAnchorAt    sql.NullTime
	ScheduleCursorAt    sql.NullTime
	NextRunAt           sql.NullTime
	CatchUpPolicy       string
	OverlapPolicy       string
	Enabled             bool
	CreatedAt           time.Time
	UpdatedAt           time.Time
}

type CandaceCronRun

type CandaceCronRun struct {
	OccurrenceID string
	JobName      string
	ScheduledAt  time.Time
	Status       string
	Attempt      int32
	WorkerID     sql.NullString
	LeaseToken   sql.NullString
	LeaseUntil   sql.NullTime
	StartedAt    sql.NullTime
	FinishedAt   sql.NullTime
	ErrorSummary sql.NullString
	SkipReason   sql.NullString
	CreatedAt    time.Time
	UpdatedAt    time.Time
}

type DBTX

type DBTX interface {
	ExecContext(context.Context, string, ...interface{}) (sql.Result, error)
	PrepareContext(context.Context, string) (*sql.Stmt, error)
	QueryContext(context.Context, string, ...interface{}) (*sql.Rows, error)
	QueryRowContext(context.Context, string, ...interface{}) *sql.Row
}

type FinishRunParams

type FinishRunParams struct {
	Status       string
	FinishedAt   sql.NullTime
	ErrorSummary sql.NullString
	OccurrenceID string
	LeaseToken   sql.NullString
}

type HasLiveRunForJobParams

type HasLiveRunForJobParams struct {
	JobName              string
	ClaimedAt            sql.NullTime
	ExcludedOccurrenceID string
}

type InsertRunningRunParams

type InsertRunningRunParams struct {
	OccurrenceID string
	JobName      string
	ScheduledAt  time.Time
	WorkerID     sql.NullString
	LeaseToken   sql.NullString
	LeaseUntil   sql.NullTime
	ClaimedAt    sql.NullTime
}

type InsertSkippedRunParams

type InsertSkippedRunParams struct {
	OccurrenceID string
	JobName      string
	ScheduledAt  time.Time
	SkippedAt    sql.NullTime
	SkipReason   sql.NullString
}

type ListExpiredRunningRunsParams

type ListExpiredRunningRunsParams struct {
	ExpiredAt sql.NullTime
	RowLimit  int32
}

type MarkExpiredRunSkippedParams

type MarkExpiredRunSkippedParams struct {
	SkippedAt    sql.NullTime
	SkipReason   sql.NullString
	OccurrenceID string
}

type MigrationSource

type MigrationSource struct {
	Files     embed.FS
	Directory string
}

MigrationSource is the embedded relational schema owned by the cron PostgreSQL adapter. An application can feed it to its migration runner without copying cron's table definitions into the application.

func EmbeddedMigrations

func EmbeddedMigrations() MigrationSource

EmbeddedMigrations returns cron's canonical PostgreSQL schema source.

type Queries

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

func New

func New(db DBTX) *Queries

func (*Queries) AcquireExpiredRun

func (q *Queries) AcquireExpiredRun(ctx context.Context, arg AcquireExpiredRunParams) (CandaceCronRun, error)

func (*Queries) AdvanceJobNextRun

func (q *Queries) AdvanceJobNextRun(ctx context.Context, arg AdvanceJobNextRunParams) (int64, error)

func (*Queries) FinishRun

func (q *Queries) FinishRun(ctx context.Context, arg FinishRunParams) (int64, error)

func (*Queries) GetJobForUpdate

func (q *Queries) GetJobForUpdate(ctx context.Context, jobName string) (CandaceCronJob, error)

func (*Queries) GetRunForUpdate

func (q *Queries) GetRunForUpdate(ctx context.Context, occurrenceID string) (CandaceCronRun, error)

func (*Queries) HasLiveRunForJob

func (q *Queries) HasLiveRunForJob(ctx context.Context, arg HasLiveRunForJobParams) (bool, error)

func (*Queries) InsertRunningRun

func (q *Queries) InsertRunningRun(ctx context.Context, arg InsertRunningRunParams) (CandaceCronRun, error)

func (*Queries) InsertSkippedRun

func (q *Queries) InsertSkippedRun(ctx context.Context, arg InsertSkippedRunParams) (CandaceCronRun, error)

func (*Queries) ListExpiredRunningRuns

func (q *Queries) ListExpiredRunningRuns(ctx context.Context, arg ListExpiredRunningRunsParams) ([]CandaceCronRun, error)

func (*Queries) ListJobs

func (q *Queries) ListJobs(ctx context.Context) ([]CandaceCronJob, error)

func (*Queries) ListJobsForUpdate

func (q *Queries) ListJobsForUpdate(ctx context.Context) ([]CandaceCronJob, error)

func (*Queries) ListRecentRuns

func (q *Queries) ListRecentRuns(ctx context.Context, rowLimit int32) ([]CandaceCronRun, error)

func (*Queries) LockReconciliation

func (q *Queries) LockReconciliation(ctx context.Context) error

func (*Queries) MarkExpiredRunSkipped

func (q *Queries) MarkExpiredRunSkipped(ctx context.Context, arg MarkExpiredRunSkippedParams) (CandaceCronRun, error)

func (*Queries) RenewRunLease

func (q *Queries) RenewRunLease(ctx context.Context, arg RenewRunLeaseParams) (int64, error)

func (*Queries) SetJobEnabled

func (q *Queries) SetJobEnabled(ctx context.Context, arg SetJobEnabledParams) (CandaceCronJob, error)

func (*Queries) UpsertJob

func (q *Queries) UpsertJob(ctx context.Context, arg UpsertJobParams) (CandaceCronJob, error)

func (*Queries) WithTx

func (q *Queries) WithTx(tx *sql.Tx) *Queries

type RenewRunLeaseParams

type RenewRunLeaseParams struct {
	LeaseUntil   sql.NullTime
	RenewedAt    time.Time
	OccurrenceID string
	LeaseToken   sql.NullString
}

type SetJobEnabledParams

type SetJobEnabledParams struct {
	Enabled bool
	JobName string
}

type Store

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

Store persists cron state in PostgreSQL. Callers own the database lifecycle.

func NewStore

func NewStore(db *sql.DB) (*Store, error)

NewStore binds the durable cron adapter to a caller-owned database pool.

func (*Store) Claim

func (*Store) Complete

func (s *Store) Complete(ctx context.Context, r cron.Completion) error

func (*Store) Expired

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

Expired returns a bounded, deterministic recovery work list without moving any cursor. Claim fences and reclaims the selected occurrence.

func (*Store) Reconcile

func (s *Store) Reconcile(ctx context.Context, defs []cron.JobDefinition, now time.Time) ([]cron.JobState, error)

func (*Store) Renew

func (s *Store) Renew(ctx context.Context, r cron.LeaseRenewal) error

func (*Store) Skip

func (s *Store) Skip(ctx context.Context, r cron.SkipRequest) error

func (*Store) Snapshot

func (s *Store) Snapshot(ctx context.Context) (cron.StoreSnapshot, error)

type UpsertJobParams

type UpsertJobParams struct {
	JobName             string
	ScheduleKind        string
	LocalHour           sql.NullInt16
	LocalMinute         sql.NullInt16
	Weekday             sql.NullInt16
	MonthDay            sql.NullInt16
	IntervalNanoseconds sql.NullInt64
	RawExpression       sql.NullString
	Timezone            string
	IntervalAnchorAt    sql.NullTime
	ScheduleCursorAt    sql.NullTime
	NextRunAt           sql.NullTime
	CatchUpPolicy       string
	OverlapPolicy       string
	Enabled             bool
}

Jump to

Keyboard shortcuts

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