pgcron

package
v0.7.0 Latest Latest
Warning

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

Go to latest
Published: Sep 29, 2026 License: MIT Imports: 14 Imported by: 0

Documentation

Overview

Package pgcron watches pg_cron jobs, which run inside Postgres where nothing can wrap them (sources/pgcron.ts). As a source, on every check it reads cron.job and declares each job with its schedule, then copies new rows of cron.job_run_details in as runs (ids "pgcron:<runid>"), so the usual evaluation raises missed, failed, stuck and slow alerts.

source := pgcron.New(db, pgcron.Options{Prefix: "db:"}) // the app's *sql.DB, any Postgres driver
cw, err := cronwatch.New(cronwatch.WithStore(store), cronwatch.WithSources(source))
cw.Start(time.Minute)

A job that is renamed, unscheduled or no longer picked keeps its old name's runs and history, and that name is declared again without a schedule, so it is never reported missed. Its description says why.

It reads through the app's database/sql handle with whatever driver the app uses (pgx's stdlib, lib/pq): a *sql.DB, a *sql.Conn or a *sql.Tx. Settings are read from pg_settings, which answers no row for a setting the role may not read where current_setting() would raise an error and abort the caller's transaction, and the source never commits or rolls back anything.

Index

Examples

Constants

View Source
const Hold = 10 * time.Minute

Hold is how long a run pg_cron has queued but not started (no start_time yet) is waited for. After that it is copied as running from when it was first seen, so a run that never starts is marked stuck like any other.

Variables

This section is empty.

Functions

func JobName

func JobName(j Job) string

JobName is the default CronWatch name for a pg_cron job, before the prefix.

func RunOf

func RunOf(row Row, job, idPrefix string, fallbackAt int64) *cronwatch.Run

RunOf is a row of cron.job_run_details as a CronWatch run, or nil for one that has not started (no start_time, not finished). A finished row with no start_time (pg_cron writes these for runs a server restart cut off, "server restarted") starts at its end_time, else at fallbackAt (the reader passes the job's newest run's start, or now).

func Schedule

func Schedule(schedule string) (string, bool)

Schedule is a pg_cron schedule as a CronWatch one: a cron expression, "$" for the last day of the month read as "L", or "N seconds" as "every Ns". pg_cron reads only the first five fields of an expression and ignores the rest, so only those are kept (a sixth would otherwise be read as seconds). It answers false for one that has no cadence to watch (@reboot).

Types

type Job

type Job struct {
	JobID int64
	// JobName is nil for a job scheduled without a name.
	JobName  *string
	Schedule string
	Database string
	Username string
	Active   bool
}

Job is a row of cron.job.

type Options

type Options struct {
	// Jobs and JobIDs pick the jobs to watch by name or id; Pick picks
	// them with a function. Default every job the role can see.
	Jobs   []string
	JobIDs []int64
	Pick   func(Job) bool
	// Prefix goes before every job name, to keep them apart from your own
	// ("db:"). It also keeps run ids apart.
	Prefix string
	// JobName is the CronWatch name for a job. Default its jobname with
	// anything other than letters, digits, ".", "_", ":" and "-" turned
	// into "-", or "pg_cron:<jobid>" when it has none (JobName). The prefix
	// goes in front either way.
	JobName func(Job) string
	// Options are job options (Grace, Timeout, MaxDuration, Expect and the
	// rest) for every job; OptionsFor gives them per job. The schedule and
	// timezone always come from pg_cron.
	Options    []cronwatch.JobOption
	OptionsFor func(Job) []cronwatch.JobOption
	// Timezone is the zone pg_cron reads its cron expressions in. Default
	// the server's cron.timezone, read from pg_settings, which shows it only
	// to roles with pg_read_all_settings; UTC (pg_cron's default) is
	// assumed when it cannot be read.
	Timezone string
}

Options configure the source.

type Querier

type Querier interface {
	QueryContext(ctx context.Context, query string, args ...any) (*sql.Rows, error)
}

Querier is what the source reads through: *sql.DB, *sql.Conn and *sql.Tx are all one.

type Row

type Row struct {
	RunID         int64
	JobID         int64
	Status        string
	ReturnMessage *string
	StartTime     *time.Time
	EndTime       *time.Time
}

Row is a row of cron.job_run_details.

type Source

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

Source is the pg_cron source. Make one with New.

func New

func New(db Querier, o Options) *Source

New is a pg_cron source reading through db.

Example

The README's pg_cron source: each check reads pg_cron's jobs and runs through the app's *sql.DB, whatever Postgres driver opened it.

package main

import (
	"database/sql"
	"os"
	"time"

	cronwatch "cronwatch.dev/go"
	"cronwatch.dev/go/pgcron"
	"cronwatch.dev/go/sqlstore"
)

func main() {
	db, err := sql.Open("pgx", os.Getenv("DATABASE_URL")) // github.com/jackc/pgx/v5/stdlib
	if err != nil {
		panic(err)
	}
	store, err := sqlstore.New(db, sqlstore.Postgres)
	if err != nil {
		panic(err)
	}
	cw, err := cronwatch.New(cronwatch.WithStore(store), cronwatch.WithSources(pgcron.New(db, pgcron.Options{Prefix: "db:"})))
	if err != nil {
		panic(err)
	}
	cw.Start(time.Minute)
	defer cw.Close()
}

func (*Source) Name

func (s *Source) Name() string

Name is "pg_cron".

func (*Source) Sync

func (s *Source) Sync(ctx context.Context, host cronwatch.SourceHost) ([]cronwatch.Alert, error)

Sync declares the jobs and copies their new runs in, returning the alerts recording them sent.

Jump to

Keyboard shortcuts

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