cronwatch

package module
v0.6.1 Latest Latest
Warning

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

Go to latest
Published: Sep 28, 2026 License: MIT Imports: 33 Imported by: 0

README

cronwatch.dev/go

Cron and scheduled-job monitoring that lives inside your Go app. Wrap a job once; every run is recorded in a database you already have, and you are told when a run is missed, fails, gets stuck, runs slow or goes over budget. No server to run, no account to make.

This is the Go port of @cronwatch/sdk, under way: the same rules, the same alert text and the same stored rows, so a Go process and a Node process can share one database, and every port reads the tables the others write. It has the core (jobs, runs, runs that span calls, checks, silences, sources, deferred delivery and the triage hook), the memory store, a database/sql store for SQLite, Postgres and MySQL, the SDK's alert channels, Claude triage, the pg_cron source, the dashboard with its JSON API, job handlers for platform crons, and integrations for robfig/cron, gocron, River and Asynq (DESIGN.md has how each part works). It is not released yet.

Docs: cronwatch.dev

Install

Go 1.25 or newer. The module requires nothing: cron expressions are read by a port of croner (the parser the SDK uses) and zones come from Go's own time package. The SQL store works over the *sql.DB your app already has, with the driver it already uses; the module imports none.

Use

package main

import (
	"context"
	"database/sql"
	"log"
	"time"

	cronwatch "cronwatch.dev/go"
	"cronwatch.dev/go/sqlstore"
	_ "modernc.org/sqlite"
)

func main() {
	db, err := sql.Open("sqlite", "file:/var/lib/app/app.db")
	if err != nil {
		log.Fatal(err)
	}
	store, err := sqlstore.New(db, sqlstore.SQLite) // or sqlstore.Postgres, sqlstore.MySQL
	if err != nil {
		log.Fatal(err)
	}
	cw := cronwatch.MustNew(
		cronwatch.WithStore(store),
		cronwatch.WithAlerts(cronwatch.ChannelFunc("pager", func(ctx context.Context, a cronwatch.Alert) error {
			return page(ctx, a.Title, a.Message)
		})),
	)
	defer cw.Close()

	nightly := cw.MustJob("nightly-report",
		cronwatch.Schedule("0 2 * * *"), cronwatch.Timezone("UTC"),
		cronwatch.Grace("15m"), cronwatch.Timeout(30*time.Minute),
		cronwatch.Expect("Report written"), cronwatch.Budget("cost", 2))

	err = nightly.Run(context.Background(), func(ctx context.Context, job *cronwatch.JobContext) error {
		path, err := buildReport(ctx) // ctx is cancelled when the job's timeout passes
		job.Log("Report written:", path) // kept with the run, shown in alerts
		job.Metric("cost", 1.2)          // watched against budgets and baselines
		return err
	})
	if err != nil {
		log.Fatal(err)
	}
}

A run is recorded when the function returns: a returned error is the failure and is returned from Run, and a panic is recorded as a failed run and carries on up the stack (a job handler answers it 500 instead, as below). The store failing never stops a job; store errors go to the error handler (cronwatch.WithErrorHandler, standard error by default). cronwatch.RunValue runs a function that returns a value: a string is the output when nothing was logged, and an *http.Response of 400 or more fails the run. Inside a job, cronwatch.Current(ctx) is its JobContext.

Schedulers

A scheduler you already run is watched with one line, each integration a module of its own so you pull only the one you use. Its jobs become CronWatch jobs with their schedules (checked against the scheduler's own fire times; one that cannot match exactly, such as a time daylight saving skips, is reported and watched without a schedule), every run is recorded, and a job taken out of the scheduler loses its schedule rather than being reported missed. Jobs are tagged with the integration and your app's name (CRONWATCH_APP_ID, else the executable's name), so two apps sharing a store never touch each other's jobs.

robfig/cron v3 (go get cronwatch.dev/go/robfigcron): a cron.Option. A job is named after its function (jobs.NightlyReport) or type; name a closure, or give options, with robfigcron.Named. watcher.Func makes a job that gets the run's context and logs.

c := cron.New(robfigcron.Watch(cw, robfigcron.Options{Chain: []cron.JobWrapper{cron.Recover(logger)}}))
c.AddFunc("0 2 * * *", jobs.NightlyReport)
c.AddJob("*/15 * * * *", robfigcron.Named("sync-invoices", syncJob, cronwatch.Grace("5m")))

gocron v2 (go get cronwatch.dev/go/gocron, gocron 2.21 or newer): its event listeners, as a scheduler option. Cron, duration, and daily, weekly and monthly jobs are read as schedules; a job is named by gocron.WithName, else its function.

s, err := gocron.NewScheduler(cwgocron.Watch(cw, cwgocron.Options{}))
s.NewJob(gocron.DailyJob(1, gocron.NewAtTimes(gocron.NewAtTime(2, 0, 0))), gocron.NewTask(jobs.NightlyReport))

River (go get cronwatch.dev/go/river, River 0.44.1 or newer): w.PeriodicJob in place of river.NewPeriodicJob, and a worker middleware that records each attempt, with cronwatch.Current(ctx) inside the worker.

w := cwriver.New(cw, cwriver.Options{Kinds: map[string][]cronwatch.JobOption{"send_invoice": nil}})
river.AddWorker(workers, w.CheckWorker())
nightly, err := cron.ParseStandard("0 2 * * *") // github.com/robfig/cron/v3, as River's docs use
config := &river.Config{
	Workers:    workers,
	Middleware: []rivertype.Middleware{w.Middleware()},
	PeriodicJobs: []*river.PeriodicJob{
		w.PeriodicJob(nightly, func() (river.JobArgs, *river.InsertOpts) {
			return NightlyReportArgs{}, nil
		}, &river.PeriodicJobOpts{ID: "nightly-report"}, cronwatch.Grace("15m")),
		w.CheckPeriodicJob(5 * time.Minute),
	},
}

Asynq (go get cronwatch.dev/go/asynq, Asynq 0.25.1 or newer): a scheduler (or periodic task manager) that declares its entries, and a server middleware.

w := cwasynq.New(cw, cwasynq.Options{})
scheduler := w.NewScheduler(redisOpt, &asynq.SchedulerOpts{Location: time.UTC})
scheduler.Register("0 2 * * *", asynq.NewTask("report:nightly", nil))
scheduler.Register("*/5 * * * *", cwasynq.CheckTask())

mux := asynq.NewServeMux()
mux.Use(w.Middleware())
mux.Handle(cwasynq.CheckType, w.CheckHandler())

Retries follow one rule everywhere: each attempt is a run, failing attempts open one alert and the attempt that succeeds closes it. An attempt given back without failing (a River snooze or cancel, an Asynq revoke) leaves no run.

Alerts

Alerts go to the console until you give channels. cronwatch.dev/go/alerts has the SDK's, on net/http alone: Slack, Discord, a signed webhook, email through Resend, Postmark, SendGrid, Mailgun or SES, SMS through Twilio, and Sentry, Honeybadger, Datadog, Rollbar, Bugsnag and New Relic.

slack, err := alerts.Slack(alerts.SlackOptions{WebhookURL: os.Getenv("SLACK_WEBHOOK_URL")})
email, err := alerts.Resend(alerts.ResendOptions{
	APIKey:       os.Getenv("RESEND_API_KEY"),
	EmailOptions: alerts.EmailOptions{From: "CronWatch <alerts@example.com>", To: []string{"ops@example.com"}},
})
cw, err := cronwatch.New(cronwatch.WithStore(store), cronwatch.WithAlerts(slack, email))

Each request is the SDK's, byte for byte, with one ten second deadline, no redirect followed, TLS verified, and errors that name only the provider and the URL's origin, with your keys cut out. Every options struct takes an HTTPClient for a proxy or a test.

Triage

cronwatch.dev/go/triage asks Claude for a short diagnosis of each alert (never per run), over plain HTTP, with no SDK to install:

diagnose, err := triage.Anthropic(triage.AnthropicOptions{Context: "A Go service on Fly.io with a Postgres database."}) // ANTHROPIC_API_KEY
cw, err := cronwatch.New(cronwatch.WithAlerts(slack), cronwatch.WithTriage(diagnose))

pg_cron

cronwatch.dev/go/pgcron watches pg_cron's jobs, which run inside Postgres where nothing can wrap them: each check reads cron.job and cron.job_run_details through your *sql.DB (any driver) and records their runs, so missed, failed, stuck and slow jobs alert like your own.

cw, err := cronwatch.New(cronwatch.WithStore(store), cronwatch.WithSources(pgcron.New(db, pgcron.Options{Prefix: "db:"})))
cw.Start(time.Minute)

Dashboard

cw.Routes(...) is the dashboard and its small JSON API as an http.Handler: every job's health, its last day and week drawn as timelines, its runs with their output, and buttons to check, silence and forget. It is the SDK's, page for page and byte for byte, so the @cronwatch/mcp server works against it as it does against a Node app. It installs as an app on a phone (a manifest, icons and a service worker, served without the token).

routes, err := cw.Routes(cronwatch.WithToken(os.Getenv("CRONWATCH_TOKEN")))
mux.Handle("/cronwatch/", routes)                               // mounted at /cronwatch
mux.Handle("/ops/cron/", http.StripPrefix("/ops/cron", routes)) // or anywhere, stripped

Everything needs the token: send it as Authorization: Bearer <token>, or open the dashboard once with ?token=<token> and a cookie keeps the browser signed in. /api/check also takes the client's cron secret, so a platform cron can run checks. With no token, the routes answer 503, except in development (CRONWATCH_ENV, APP_ENV or GO_ENV set to development, dev, local, test or testing), where they make one and print a sign-in link on the first request; cronwatch.WithoutToken() serves them open, behind your own auth. The base path is found from the request: what http.StripPrefix took off, else the part of the ServeMux pattern before its wildcard or trailing slash, else /cronwatch; cronwatch.WithBasePath sets it. Behind a proxy, cronwatch.WithOrigin("https://app.example.com") or cronwatch.WithTrustProxy() gives the public origin the cross-site check and the cookie use.

Platform crons

job.Handler(fn) is a job as an http.Handler, for a cron that calls a URL (Cloud Scheduler, Vercel, a crontab line running curl). A request must carry Authorization: Bearer <CRON_SECRET>; each one runs the function as a recorded run and is answered with how it went.

mux.Handle("POST /api/cron/nightly", nightly.Handler(func(ctx context.Context, job *cronwatch.JobContext, w http.ResponseWriter, r *http.Request) error {
	return buildReport(ctx)
}))

The answer is {"ok","job","run","status","durationMs"}, 200 or 500 (a panic in the function too, recorded as a failed run), unless the function wrote its own, in which case a status of 400 or more fails the run. cronwatch.HandlerValue takes a function returning a value, such as the *http.Response of a call it made, which becomes the answer. Without a secret the handler answers 503 outside development; cronwatch.WithoutSecret() lets anyone in. cronwatch.Lambda(handler) is any handler (a job's, or the dashboard) as an AWS Lambda function behind API Gateway or a function URL, for lambda.Start, with no AWS module in this one.

Checks

Missed and stuck runs are found by a check. A long-running service (a server, a worker, a process running a scheduler) calls cw.Start(time.Minute), a goroutine that checks every minute until cw.Stop or cw.Close; one process is enough, and more are harmless. A program run from a crontab checks from a second crontab line, on a store both reach (examples/crontab is one):

# m  h  dom mon dow  command
0    2  *   *   *    /usr/local/bin/nightly report
*/5  *  *   *   *    /usr/local/bin/nightly check

Both commands declare the job, so the check knows its schedule before its first run; report runs it with job.Run and exits non-zero on the error Run returns, and check calls cw.Check(ctx).

A run that starts in one call and ends in another (work handed to a queue, a webhook that reports back later) is one run too:

h, err := nightly.Start(ctx, cronwatch.WithRunID(deliveryID))
// ... later, perhaps in another process:
h, err = nightly.Resume(ctx, deliveryID)
h.Log("done")
h.Finish(ctx) // or h.Fail(ctx, err); a run is judged once, however many processes finish it

Testing

cd packages/go && go test -race ./...           # the core, the dashboard, the channels, triage and pg_cron (against fakes), standard library only
cd packages/go/sqltest && go test -race ./...   # the SQL store with real drivers
cd packages/go/robfigcron && go test -race ./... # and gocron, river, asynq, examples: each a module of its own

The integrations require their schedulers at the oldest release they support; CI also tests the newest (go get it first). River's tests need a Postgres (CRONWATCH_TEST_PG, each run in a schema of its own) and Asynq's end-to-end test a Redis (CRONWATCH_TEST_REDIS=redis://127.0.0.1:56379/0, from docker run -d --name cw-redis -p 56379:6379 redis:7-alpine); both skip without them.

Run npm ci && npm run build at the repository root first: the croner parity test (internal/schedule) and the SQLite file shared with Node (sqltest) use the built SDK, and skip, with the reason, without it. The dashboard is checked against packages/ruby/test/web/golden.json, the SDK's answers to a fixed seed, and CRONWATCH_TEST_GO=1 npm test --workspace packages/mcp drives the MCP server against it. The tests set the process zone to UTC, as the conformance fixtures are made in UTC.

sqltest is a module of its own so that the drivers it uses (modernc.org/sqlite, github.com/jackc/pgx/v5, github.com/lib/pq, github.com/go-sql-driver/mysql) never appear in the cronwatch.dev/go module. SQLite always runs; Postgres, MySQL and MariaDB run when these are set, and the pg_cron source against a real pg_cron (through pgx and lib/pq) when CRONWATCH_TEST_PGCRON names a Postgres with the extension preloaded:

docker run -d --name cw-pg -e POSTGRES_PASSWORD=cw -p 55432:5432 postgres:17
docker run -d --name cw-mysql -e MYSQL_ROOT_PASSWORD=cw -e MYSQL_DATABASE=cw -p 53306:3306 mysql:8.4
docker run -d --name cw-mariadb -e MARIADB_ROOT_PASSWORD=cw -e MARIADB_DATABASE=cw -p 53307:3306 mariadb:11.4
# pg_cron: postgres:16 with postgresql-16-cron installed, run with
#   postgres -c shared_preload_libraries=pg_cron -c cron.database_name=cw

CRONWATCH_TEST_PG=postgres://postgres:cw@127.0.0.1:55432/postgres \
CRONWATCH_TEST_MYSQL=mysql://root:cw@127.0.0.1:53306/cw \
CRONWATCH_TEST_MARIADB=mysql://root:cw@127.0.0.1:53307/cw \
CRONWATCH_TEST_PGCRON=postgres://postgres:cw@127.0.0.1:55433/cw \
go test -race ./...

Documentation

Overview

Package cronwatch is cron and scheduled-job monitoring that lives inside a Go app: wrap a job once, and every run is recorded in a database the app already has, and alerts go out when a run is missed, fails, gets stuck, runs slow or goes over budget. It is the Go port of @cronwatch/sdk, with the same rules, alert text and stored rows, so a Go process can share a database with the SDK and its other ports.

A client is made once, jobs are declared on it, and each run of a job is its function called through Run:

cw, err := cronwatch.New(cronwatch.WithStore(store))
nightly, err := cw.Job("nightly-report", cronwatch.Schedule("0 2 * * *"), cronwatch.Grace("15m"))
err = nightly.Run(ctx, func(ctx context.Context, job *cronwatch.JobContext) error {
	job.Log("Report written")
	return nil
})

Missed and stuck runs are found by Check, called on an interval by Start in a long-running service, or from a crontab line. The sqlstore package keeps everything in SQLite, Postgres or MySQL over the app's own *sql.DB; MemoryStore, the default, forgets on restart. The alerts package holds the SDK's alert channels (Slack, Discord, a webhook, email, SMS and error trackers), the triage package Claude triage, and the pgcron package a source that watches pg_cron's jobs.

The dashboard and its JSON API are an http.Handler, Client.Routes, and a job run by a platform cron that calls a URL is one too, Job.Handler:

routes, err := cw.Routes(cronwatch.WithToken(os.Getenv("CRONWATCH_TOKEN")))
mux.Handle("/cronwatch/", routes)
mux.Handle("POST /api/cron/nightly", nightly.Handler(buildReport))

Every value a store holds is the SDK's JSON: each type's MarshalJSON writes it byte for byte, with keys in JavaScript's order.

Example

A job declared once and run: the failure is recorded, judged and alerted.

package main

import (
	"context"
	"errors"
	"fmt"
	"time"

	cronwatch "cronwatch.dev/go"
)

func main() {
	ctx := context.Background()
	now := time.Date(2026, 1, 5, 2, 0, 0, 0, time.UTC).UnixMilli()
	alerts := cronwatch.ChannelFunc("print", func(_ context.Context, a cronwatch.Alert) error {
		fmt.Println(a.Title)
		fmt.Println(a.Message)
		return nil
	})
	cw := cronwatch.MustNew(cronwatch.WithAlerts(alerts), cronwatch.WithClock(func() int64 { return now }))
	nightly := cw.MustJob("nightly-report", cronwatch.Schedule("0 2 * * *"), cronwatch.Timezone("UTC"), cronwatch.Grace("15m"))

	err := nightly.Run(ctx, func(ctx context.Context, job *cronwatch.JobContext) error {
		job.Log("connecting to postgres://app:" + "hunter" + "2@db.internal/app")
		return errors.New("connection refused")
	})
	fmt.Println("returned:", err)

	runs, _ := cw.Runs(ctx, "nightly-report", 1)
	fmt.Println(runs[0].Status, *runs[0].Output)
}
Output:
nightly-report failed
Started 2026-01-05 02:00:00 UTC (now), ran 0ms.
Error: connection refused
Output (tail):
connecting to postgres://app:[redacted]@db.internal/app
returned: connection refused
failed connecting to postgres://app:[redacted]@db.internal/app

Index

Examples

Constants

View Source
const (

	// DefaultBasePath is where the dashboard is taken to be mounted when
	// nothing else says: not WithBasePath, not http.StripPrefix, not a
	// ServeMux pattern.
	DefaultBasePath = "/cronwatch"
)
View Source
const MaxBody = 1 << 20

MaxBody is the most of a request body the dashboard reads: its forms and JSON are a few bytes. A body past it is answered 413; the SDK leaves this to the server in front of it.

View Source
const ReservedRunIDPrefix = "pgcron:"

ReservedRunIDPrefix starts the run ids of the pg_cron source, so no other run may use it.

View Source
const Version = "0.6.1"

Version is the release this module is, the same as every CronWatch package's. scripts/release.mjs bumps it; the module's version is its tag, packages/go/vX.Y.Z.

Variables

View Source
var Stderr io.Writer = os.Stderr

Stderr is where the default error handler, the in-memory store's warning and the console channel's failures go. Tests may replace it.

View Source
var Stdout io.Writer = os.Stdout

Stdout is where the console channel writes recoveries.

Functions

func HandlerValue

func HandlerValue[T any](job *Job, fn func(ctx context.Context, job *JobContext, w http.ResponseWriter, r *http.Request) (T, error), options ...HandlerOption) http.Handler

HandlerValue is Job.Handler for a function that returns a value, as RunValue is Run's: a string is the run's output when nothing was logged, and an *http.Response (a call to another service, say) is the answer, copied to the caller, its status of 400 or more failing the run with "HTTP <status> <reason>". The response's body is closed.

func Lambda

func Lambda(h http.Handler) func(ctx context.Context, event LambdaEvent) (LambdaResult, error)

Lambda is h as an AWS Lambda function behind API Gateway or a function URL, for aws-lambda-go's lambda.Start, which the app imports:

lambda.Start(cronwatch.Lambda(nightly.Handler(buildReport)))

The event becomes an *http.Request over https (the host from its Host header or its domain name, the body decoded when it is base64), and h's answer the proxy result, its body as text when it is UTF-8 and base64 otherwise, without a Content-Length (API Gateway sets its own). A directly invoked function (EventBridge Scheduler) has no headers to carry a bearer; IAM decides who may invoke it, so such a handler takes WithoutSecret.

func RedactSecrets

func RedactSecrets(text string) string

RedactSecrets is the default redaction: it blanks values that look like secrets (secret-named pairs, credentials in URLs, authorization headers, private keys, JWTs, webhook URLs, and common API key formats), exactly as the SDK's default does. A WithRedact function can call it and add its own patterns on top.

Example

A WithRedact function that keeps the default and adds a pattern of its own.

card := regexp.MustCompile(`\d{16}`)
redact := func(text string) string { return card.ReplaceAllString(RedactSecrets(text), "[card]") }
fmt.Println(redact("charged 4242424242424242"))
Output:
charged [card]

func RunValue

func RunValue[T any](ctx context.Context, job *Job, fn func(ctx context.Context, job *JobContext) (T, error), options ...RunOption) (T, error)

RunValue runs fn as a recorded run of job, as Job.Run does, and returns what it returns. A string is the run's output when nothing was logged (and what an expect rule checks), and an *http.Response whose status is 400 or more fails the run with "HTTP <status> <reason>".

Types

type Alert

type Alert struct {
	Type AlertType
	// The run behind the alert, when there is one.
	Run     *Run
	Details AlertDetails
	Job     string
	// The job's definition when the alert was made.
	Definition Definition
	// One line, suitable as a notification title.
	Title string
	// A few lines of plain text with the specifics.
	Message string
	// A short diagnosis from the triage function, when one is configured.
	// Nil with TriageTried set means triage was tried and gave nothing; it
	// is not tried again for this alert.
	Triage      *string
	TriageTried bool
	At          int64
}

Alert is a condition opening or closing, with the text every channel shows.

func (Alert) JSValue

func (a Alert) JSValue() any

JSValue is the alert as the SDK writes it.

func (Alert) MarshalJSON

func (a Alert) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

func (*Alert) UnmarshalJSON

func (a *Alert) UnmarshalJSON(b []byte) error

UnmarshalJSON reads the SDK's JSON.

type AlertDetails

type AlertDetails interface {
	// contains filtered or unexported methods
}

AlertDetails is what an alert carries beyond its title and message: a MissedDetails, FailureDetails (failed and stuck), SlowDetails, OverBudgetDetails or RecoveredDetails.

type AlertType

type AlertType string

AlertType is a condition opening, or "recovered".

const (
	AlertMissed     AlertType = "missed"
	AlertFailed     AlertType = "failed"
	AlertStuck      AlertType = "stuck"
	AlertSlow       AlertType = "slow"
	AlertOverBudget AlertType = "over_budget"
	AlertRecovered  AlertType = "recovered"
)

The alert types.

type BudgetBreach

type BudgetBreach struct {
	Metric string
	Value  float64
	Limit  float64
	// "budget", or how the baseline was worked out.
	Basis string
}

BudgetBreach is one metric over its ceiling or its baseline.

type Channel

type Channel interface {
	// Name names the channel in errors ("alert channel <name>").
	Name() string
	// Send returns once the alert went out (to at least one recipient), or
	// an error when it went nowhere. ctx ends after 15 seconds.
	Send(ctx context.Context, alert Alert, cc ChannelContext) error
}

Channel is where alerts go.

func ChannelFunc

func ChannelFunc(name string, send func(ctx context.Context, alert Alert) error) Channel

ChannelFunc wraps a function as a channel (the SDK's custom()).

func Console

func Console() Channel

Console is the default channel: it writes each alert to standard error, and recoveries to standard output.

type ChannelContext

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

ChannelContext is what the client hands a channel with each alert.

func NewChannelContext

func NewChannelContext(report func(error)) ChannelContext

NewChannelContext is a channel context that reports to report, for sending to a channel outside a client (a test, say).

func (ChannelContext) ReportError

func (cc ChannelContext) ReportError(err error)

ReportError reports a problem that did not stop the alert going out, such as one of several recipients refusing it. It goes to the client's error handler; the zero ChannelContext writes it to standard error, as the SDK does for a channel called without one.

type CheckResult

type CheckResult struct {
	CheckedAt int64
	Jobs      []JobSummary
	Alerts    []Alert
	Pruned    int
}

CheckResult is what a check found and sent.

func (CheckResult) JSValue

func (c CheckResult) JSValue() any

JSValue is the result as the SDK writes it.

func (CheckResult) MarshalJSON

func (c CheckResult) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

type Client

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

Client watches an app's jobs: it records their runs in a store, judges each one, sends alerts, and runs the checks that find missed and stuck runs. One per app, made once with New. It is safe for use by many goroutines at once.

func MustNew

func MustNew(options ...Option) *Client

MustNew is New for a package-level client: it panics when an option is invalid.

func New

func New(options ...Option) (*Client, error)

New makes a client. With no options it keeps everything in memory and writes alerts to the console.

func (*Client) Check

func (c *Client) Check(ctx context.Context) (*CheckResult, error)

Check looks for missed and stuck runs across every job, sends alerts, retries alerts no channel accepted, and prunes old runs. Call it from an interval (Start), a cron hitting the dashboard's check route, or by hand. Concurrent calls share one check. A job that cannot be evaluated is reported to the error handler and shown as failing; the error returned is for the store failing as the check starts.

The shared check runs to the end without the first caller's cancellation (its values still reach the store), so a caller that gives up (a request that ended) neither fails the check for the others nor leaves a job half checked; each caller stops waiting when its own ctx ends, with its cause.

func (*Client) Close

func (c *Client) Close() error

Close stops the interval and closes the store.

func (*Client) CronSecret

func (c *Client) CronSecret() string

CronSecret is the secret job handlers' requests must carry, or "" when none is set.

func (*Client) DefinedJobs

func (c *Client) DefinedJobs() []Definition

DefinedJobs are the definitions declared in this process, in the order they were first declared.

func (*Client) Forget

func (c *Client) Forget(ctx context.Context, name string) error

Forget removes a job and its runs from the store. A job still declared in code comes back on its next run.

func (*Client) GetRun

func (c *Client) GetRun(ctx context.Context, id string) (*Run, error)

GetRun is one run by its id, or nil.

func (*Client) Job

func (c *Client) Job(name string, options ...JobOption) (*Job, error)

Job declares a job and returns its handle. Call it once, at startup, and keep the handle. Declaring a name again replaces its definition.

Example

Options keep the order they are given in, as the SDK's object literal does, so the stored definition is the same JSON a Node process writes.

package main

import (
	"fmt"
	"time"

	cronwatch "cronwatch.dev/go"
)

func main() {
	cw := cronwatch.MustNew()
	job := cw.MustJob("sync", cronwatch.Schedule("every 5m"), cronwatch.Timeout(90*time.Second),
		cronwatch.Budget("tokens", 5000), cronwatch.Expect("synced"))
	out, _ := job.Definition().MarshalJSON()
	fmt.Println(string(out))
}
Output:
{"schedule":"every 5m","timeout":90000,"budget":{"tokens":5000},"name":"sync","expect":"contains \"synced\""}

func (*Client) JobSummary

func (c *Client) JobSummary(ctx context.Context, name string) (*JobSummary, error)

JobSummary is one job's summary, or nil when the store does not know it.

func (*Client) Jobs

func (c *Client) Jobs(ctx context.Context) ([]JobSummary, error)

Jobs is every job the store knows about, with its health. It sends no alerts.

func (*Client) JobsWithRuns

func (c *Client) JobsWithRuns(ctx context.Context, limit int) ([]JobWithRuns, error)

JobsWithRuns is every job's summary with its newest limit runs (0 to 500), read together. What the dashboard shows.

func (*Client) MustJob

func (c *Client) MustJob(name string, options ...JobOption) *Job

MustJob is Job for package-level declarations: it panics when an option is invalid.

func (*Client) MustRoutes

func (c *Client) MustRoutes(options ...RoutesOption) *Routes

MustRoutes is Routes for a package-level declaration: it panics where Routes returns an error.

func (*Client) Now

func (c *Client) Now() int64

Now is the client's clock, in epoch milliseconds.

func (*Client) RecordRun

func (c *Client) RecordRun(ctx context.Context, input Run, options ...RecordOption) ([]Alert, error)

RecordRun records a run that happened outside this process, for a Source. Its job must be declared with Job first. Runs are keyed by id: a new one is inserted, a stored one still running (or marked timeout by a check) is finished when this one is not running, and anything else is left alone, so recording the same run twice changes nothing. Finishing is conditional: when two processes record the same finish, only the one whose write lands evaluates it, and the other reports it as already finished. A stored run of another job is left alone and reported. A finished run is judged as if it had been wrapped here (expect, failures, duration, budgets) and its output and error are redacted the same way. Returns the alerts it sent.

func (*Client) ReportError

func (c *Client) ReportError(err error, where string)

ReportError hands an error to the client's error handler, as the client reports its own. A source uses it.

func (*Client) ResumeRun

func (c *Client) ResumeRun(ctx context.Context, name, runID string) (*RunHandle, error)

ResumeRun is Resume for a job declared in this process, by name.

func (*Client) Routes

func (c *Client) Routes(options ...RoutesOption) (*Routes, error)

Routes is the dashboard and its JSON API as an http.Handler, the SDK's cw.routes(). Mount it under a prefix with http.StripPrefix or a ServeMux pattern:

routes, err := cw.Routes(cronwatch.WithToken(os.Getenv("CRONWATCH_TOKEN")))
mux.Handle("/cronwatch/", routes)                              // the base is /cronwatch
mux.Handle("/ops/cron/", http.StripPrefix("/ops/cron", routes)) // the base is /ops/cron

The error is for a WithOrigin value that is not an http or https URL.

Example

The dashboard and its JSON API, mounted under a prefix: the base path its links use is found from the mount.

package main

import (
	"fmt"
	"net/http"
	"net/http/httptest"

	cronwatch "cronwatch.dev/go"
)

func main() {
	cw := cronwatch.MustNew(cronwatch.WithoutCronSecret())
	routes, err := cw.Routes(cronwatch.WithToken("a-long-random-token"))
	if err != nil {
		panic(err)
	}
	mux := http.NewServeMux()
	mux.Handle("/ops/cron/", http.StripPrefix("/ops/cron", routes))

	req := httptest.NewRequest("GET", "/ops/cron/api/jobs", nil)
	req.Header.Set("Authorization", "Bearer a-long-random-token")
	rec := httptest.NewRecorder()
	mux.ServeHTTP(rec, req)
	fmt.Println(rec.Code, rec.Body.String())
}
Output:
200 {"ok":true,"jobs":[]}

func (*Client) Run

func (c *Client) Run(ctx context.Context, name string, fn JobFunc, options ...JobOption) error

Run runs a job by name without keeping a handle, declaring it on first use (or again, when options are given).

func (*Client) Runs

func (c *Client) Runs(ctx context.Context, name string, limit int) ([]Run, error)

Runs is a job's runs, newest first. limit is 1 to 500.

func (*Client) Silence

func (c *Client) Silence(ctx context.Context, name string, d time.Duration) (JobState, error)

Silence stops alerts for a job for a while. State keeps updating underneath: nothing opens while it is silenced, so the first problem after the silence alerts as usual.

func (*Client) Start

func (c *Client) Start(every time.Duration)

Start checks on an interval in a goroutine, for long-running servers: the first check a second from now, then every `every` (a minute when 0, five seconds at least). Not for serverless functions, where nothing runs between requests: call Check from a cron there instead. A second Start does nothing.

func (*Client) Stop

func (c *Client) Stop()

Stop stops the interval Start began. A check in flight finishes.

func (*Client) Store

func (c *Client) Store() Store

Store is where this client keeps jobs, runs and state.

func (*Client) SyncJob

func (c *Client) SyncJob(ctx context.Context, name string) (bool, error)

SyncJob writes the definition declared in this process under name to the store now, unless the store already holds that definition (whatever order its keys come back in), and says whether it wrote. A run or a check writes a declaration anyway, once; a scheduler integration calls SyncJob so a process that only schedules (and never runs or checks) still puts its jobs where the processes that run and check them read them, and so a definition another process changed since (took the schedule out of, say) is written back. A name not declared here is an error.

func (*Client) Unsilence

func (c *Client) Unsilence(ctx context.Context, name string) (JobState, error)

Unsilence ends a silence.

type Condition

type Condition string

Condition is something wrong with a job that opens once, alerts, and closes with a recovery.

const (
	ConditionMissed     Condition = "missed"
	ConditionFailed     Condition = "failed"
	ConditionStuck      Condition = "stuck"
	ConditionSlow       Condition = "slow"
	ConditionOverBudget Condition = "over_budget"
)

The conditions, in the SDK's order.

type Definition

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

Definition is a job's definition as a store holds it: the SDK's JSON object, its fields in the order they were given (defaults, then the job's options, then name), `expect` described in words. Fields a newer writer added are kept.

func DescribeJob

func DescribeJob(name string, options ...JobOption) Definition

DescribeJob is the definition these options give a job, before any client's defaults and without checking them: what a source compares to tell whether a job it declares has changed (the SDK compares the options object's JSON).

func (Definition) Description

func (d Definition) Description() string

Description is the job's description, or "".

func (Definition) Expect

func (d Definition) Expect() string

Expect describes the expect rule ("contains \"done\""), or "".

func (Definition) Get

func (d Definition) Get(key string) (any, bool)

Get is a field as JSON reads it (a string, a float64, a bool, nil, a []any or a map), and whether it is there.

func (Definition) JSValue

func (d Definition) JSValue() any

JSValue is the definition as a JSON object.

func (Definition) Keys

func (d Definition) Keys() []string

Keys are the fields present, in order.

func (Definition) MarshalJSON

func (d Definition) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

func (Definition) Name

func (d Definition) Name() string

Name is the job's name.

func (Definition) Schedule

func (d Definition) Schedule() string

Schedule is the cron expression or "every <duration>", or "".

func (Definition) Tags

func (d Definition) Tags() []string

Tags are the job's tags.

func (Definition) Timezone

func (d Definition) Timezone() string

Timezone is the IANA zone the schedule is read in, or "".

func (*Definition) UnmarshalJSON

func (d *Definition) UnmarshalJSON(b []byte) error

UnmarshalJSON reads a JSON object.

type DeliverMode

type DeliverMode string

DeliverMode says where alerts are sent from.

const (
	// DeliverNow sends each alert from the process that produced it. The default.
	DeliverNow DeliverMode = "now"
	// DeliverAtCheck sends nothing from this process: each alert is queued
	// in the store, and the next check in a process that delivers now sends
	// it (with triage). For a process that records runs but cannot reach
	// the network, such as a sandboxed backup job.
	DeliverAtCheck DeliverMode = "check"
)

type Duration

type Duration interface {
	~string | ~int | ~int64 | ~float64
}

Duration is what a duration option takes: text like "15m", "1h30m", "90s" or "2d" (the SDK's form, kept as written in the stored definition), a time.Duration, or a whole or fractional number of milliseconds. A time.Duration is stored as its milliseconds, as the SDK stores a number.

type FailureDetails

type FailureDetails struct {
	ConsecutiveFailures int
	Threshold           int
}

FailureDetails counts the failures behind a failed or stuck alert.

type HandlerFunc

type HandlerFunc func(ctx context.Context, job *JobContext, w http.ResponseWriter, r *http.Request) error

HandlerFunc is a job's work for one request. ctx is the request's context with the job's timeout added. The function may answer the request itself through w, and a status of 400 or more written there fails the run; when it writes nothing, the handler answers with how the run went.

type HandlerOption

type HandlerOption func(*handlerConfig)

HandlerOption configures a job's handler.

func WithSecret

func WithSecret(secret string) HandlerOption

WithSecret is the secret the handler's requests must carry, as `Authorization: Bearer <secret>`, in place of the client's cron secret (WithCronSecret, $CRON_SECRET by default). "" counts as unset.

func WithoutSecret

func WithoutSecret() HandlerOption

WithoutSecret lets anyone run the job through the handler (the SDK's `secret: null`), for one behind your own auth, or a function only its platform can invoke.

type Job

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

Job is a declared job's handle.

func (*Job) Definition

func (j *Job) Definition() Definition

Definition is the job's definition as it is stored.

func (*Job) Handler

func (j *Job) Handler(fn HandlerFunc, options ...HandlerOption) http.Handler

Handler is the job as an http.Handler (the SDK's job.handler()). A request must send `Authorization: Bearer <secret>` (compared in constant time): the handler's WithSecret, else the client's cron secret. With no secret at all it answers 503 and reports it once to the error handler as "handler", unless the environment is development (see WithToken) or WithoutSecret (or the client's WithoutCronSecret) lets anyone in; a wrong or missing bearer is 401. Each request it lets in runs fn as a run with the trigger "handler", answered with {"ok","job","run","status","durationMs"}, 200 when the run was ok and 500 when it failed, with the error's first line as "error" for a caller who sent the secret. A function that wrote its own answer is answered with that, and a status of 400 or more it wrote fails the run. A panic in fn is a failed run like an error (`panic: <value>`), answered 500 as the SDK answers a throw; a panic with http.ErrAbortHandler, or one after fn began its own answer, is recorded the same way and then aborts the response, as net/http does.

mux.Handle("POST /api/cron/nightly", nightly.Handler(func(ctx context.Context, job *cronwatch.JobContext, w http.ResponseWriter, r *http.Request) error {
	return buildReport(ctx)
}))
Example

A job run by a platform cron that calls a URL with the cron secret.

package main

import (
	"context"
	"encoding/json"
	"fmt"
	"net/http"
	"net/http/httptest"

	cronwatch "cronwatch.dev/go"
)

func main() {
	cw := cronwatch.MustNew(cronwatch.WithCronSecret("the-cron-secret"))
	nightly := cw.MustJob("nightly-report")
	handler := nightly.Handler(func(ctx context.Context, job *cronwatch.JobContext, w http.ResponseWriter, r *http.Request) error {
		job.Log("Report written")
		return nil
	})

	req := httptest.NewRequest("POST", "/api/cron/nightly", nil)
	req.Header.Set("Authorization", "Bearer the-cron-secret")
	rec := httptest.NewRecorder()
	handler.ServeHTTP(rec, req)
	var answer struct {
		OK     bool   `json:"ok"`
		Status string `json:"status"`
	}
	_ = json.Unmarshal(rec.Body.Bytes(), &answer)
	fmt.Println(rec.Code, answer.OK, answer.Status)
}
Output:
200 true ok

func (*Job) Name

func (j *Job) Name() string

Name is the job's name.

func (*Job) Resume

func (j *Job) Resume(ctx context.Context, runID string) (*RunHandle, error)

Resume is a handle on a run this job started elsewhere, by its id, so this process can log to it and finish it.

func (*Job) Run

func (j *Job) Run(ctx context.Context, fn JobFunc, options ...RunOption) error

Run runs fn now as a recorded run and returns its error. The run is recorded however the store is doing: store errors go to the error handler, never to the caller. A panic in fn is recorded as a failed run and then carries on up the stack.

func (*Job) Start

func (j *Job) Start(ctx context.Context, options ...RunOption) (*RunHandle, error)

Start records a running run now, to finish later with the handle, perhaps from another process (see Resume). Store failures go to the error handler; it returns an error only for an invalid run id or one that belongs to another job. A run that is never finished is marked stuck by the first check after the job's timeout.

Example

A run that starts in one call and ends in another, as the README shows: work handed to a queue, a webhook that reports back later.

package main

import (
	"context"
	"fmt"

	cronwatch "cronwatch.dev/go"
)

func main() {
	ctx := context.Background()
	cw := cronwatch.MustNew(cronwatch.WithoutCronSecret())
	nightly := cw.MustJob("nightly-report")
	h, err := nightly.Start(ctx, cronwatch.WithRunID("delivery-42"))
	if err != nil {
		panic(err)
	}
	// ... later, perhaps in another process:
	h, err = nightly.Resume(ctx, "delivery-42")
	if err != nil {
		panic(err)
	}
	h.Log("done")
	run := h.Finish(ctx)
	fmt.Println(run.ID, run.Status, *run.Output)
}
Output:
delivery-42 ok done

type JobContext

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

JobContext is what a job's function gets: its run, and where it logs output and reports metrics.

func Current

func Current(ctx context.Context) *JobContext

Current is the job context of the run ctx belongs to, or nil outside one.

func (*JobContext) Log

func (j *JobContext) Log(parts ...any)

Log appends a line of output: the parts joined by spaces, strings as they are, errors as "Name: message", anything else as JSON. Kept with the run (the last 16 KB), shown in alerts and on the dashboard.

func (*JobContext) Metric

func (j *JobContext) Metric(name string, value float64) error

Metric reports a number for this run: tokens, cost, rows, anything. Watched against budgets and baselines. A later value for the same name replaces an earlier one.

func (*JobContext) Metrics

func (j *JobContext) Metrics(values Metrics) error

Metrics reports several numbers at once.

func (*JobContext) Name

func (j *JobContext) Name() string

Name is the job's name.

func (*JobContext) RunID

func (j *JobContext) RunID() string

RunID is the run's id.

func (*JobContext) StartedAt

func (j *JobContext) StartedAt() int64

StartedAt is when the run started, in epoch milliseconds.

type JobFunc

type JobFunc func(ctx context.Context, job *JobContext) error

JobFunc is a job's work. ctx is cancelled when the job's timeout passes (and when the caller's context is); job logs output and reports metrics. A returned error fails the run and is returned from Run.

type JobHealth

type JobHealth string

JobHealth is how a job looks at a glance.

const (
	HealthHealthy  JobHealth = "healthy"
	HealthLate     JobHealth = "late"
	HealthFailing  JobHealth = "failing"
	HealthStuck    JobHealth = "stuck"
	HealthSilenced JobHealth = "silenced"
	HealthNeverRan JobHealth = "never_ran"
)

The healths a summary reports.

type JobOption

type JobOption func(*jobConfig)

JobOption configures a job: when it runs and what counts as trouble.

func Budget

func Budget(metric string, ceiling float64) JobOption

Budget sets a ceiling for a metric reported with Metric: Budget("cost", 2) alerts when a run reports cost above 2. Give it once per metric. Metrics without a ceiling alert when a run reports more than three times the recent median, once there are five runs to compare against.

func Description

func Description(text string) JobOption

Description describes the job on the dashboard.

func Expect

func Expect(text string) JobOption

Expect makes a successful run fail unless its output contains text. Catches the job that exits cleanly and did nothing.

func ExpectFunc

func ExpectFunc(fn func(output string) bool) JobOption

ExpectFunc makes a successful run fail unless fn returns true for its output. A panic in fn fails the run with what it panicked with.

func ExpectMatch

func ExpectMatch(re *regexp.Regexp) JobOption

ExpectMatch makes a successful run fail unless the regular expression matches its output.

func FailuresBeforeAlert

func FailuresBeforeAlert(n int) JobOption

FailuresBeforeAlert alerts on the nth consecutive failure rather than the first. Default 1.

func Grace

func Grace[D Duration](d D) JobOption

Grace is how late a run may start before it counts as missed. Default 10m.

func MaxDuration

func MaxDuration[D Duration](d D) JobOption

MaxDuration alerts when a successful run takes longer. Without it, a run is slow when it takes more than twice the p95 of recent runs (and over 10s), once there are five runs to compare against.

func Schedule

func Schedule(expr string) JobOption

Schedule is when the job is supposed to run: a five or six field cron expression ("0 2 * * *"), a nickname ("@hourly"), or an interval ("every 5m"). Leave it out for a job with no fixed cadence: failures, duration and budgets are still watched, but nothing is ever missed.

func Tags

func Tags(tags ...string) JobOption

Tags label the job.

func Timeout

func Timeout[D Duration](d D) JobOption

Timeout is how long a run may go on before it is treated as stuck and marked timeout. Default 1h. The context a job's function gets is cancelled when it passes.

func Timezone

func Timezone(name string) JobOption

Timezone is the IANA zone the cron expression is read in. The default is the process's zone (time.Local). Vercel and GitHub Actions run their crons in UTC.

type JobState

type JobState struct {
	Job string
	// Conditions currently open, in the order they opened.
	Open                []OpenCondition
	ConsecutiveFailures int
	SilencedUntil       *int64
	// When an alert last reached at least one channel.
	LastAlertAt *int64
	// Conditions that alerted and have since closed, waiting for the
	// recovered alert the next successful run sends. Nil is a state written
	// before the field existed (absent from the JSON).
	PendingRecovery []Condition
	// Alerts no channel accepted, each retried once per check. Nil is absent.
	Undelivered []Alert
	// Goes up by one on every write (see Store.CompareAndSetState). Nil is
	// a state written before versions, which counts as 0.
	Version *int64
	// contains filtered or unexported fields
}

JobState is what the checks remember about a job between runs.

func (JobState) JSValue

func (s JobState) JSValue() any

JSValue is the state as the SDK writes it.

func (JobState) MarshalJSON

func (s JobState) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

func (*JobState) UnmarshalJSON

func (s *JobState) UnmarshalJSON(b []byte) error

UnmarshalJSON reads the SDK's JSON.

type JobSummary

type JobSummary struct {
	Name       string
	Definition Definition
	Health     JobHealth
	Open       []Condition
	LastRun    *Run
	// When the schedule says the next run is due. Nil without a schedule.
	NextExpectedAt      *int64
	ConsecutiveFailures int
	SilencedUntil       *int64
	Stats               Stats
}

JobSummary is a job and its health, as the dashboard shows it.

func (JobSummary) JSValue

func (s JobSummary) JSValue() any

JSValue is the summary as the SDK writes it.

func (JobSummary) MarshalJSON

func (s JobSummary) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

type JobWithRuns

type JobWithRuns struct {
	Job  JobSummary
	Runs []Run
}

JobWithRuns is a job's summary and its newest runs.

type LambdaEvent

type LambdaEvent struct {
	Version string `json:"version,omitempty"`
	// 1.0
	HTTPMethod                      string              `json:"httpMethod,omitempty"`
	Path                            string              `json:"path,omitempty"`
	MultiValueHeaders               map[string][]string `json:"multiValueHeaders,omitempty"`
	QueryStringParameters           map[string]string   `json:"queryStringParameters,omitempty"`
	MultiValueQueryStringParameters map[string][]string `json:"multiValueQueryStringParameters,omitempty"`
	// 2.0
	RawPath        string   `json:"rawPath,omitempty"`
	RawQueryString string   `json:"rawQueryString,omitempty"`
	Cookies        []string `json:"cookies,omitempty"`
	// Both
	Headers         map[string]string    `json:"headers,omitempty"`
	RequestContext  LambdaRequestContext `json:"requestContext"`
	Body            string               `json:"body,omitempty"`
	IsBase64Encoded bool                 `json:"isBase64Encoded,omitempty"`
}

LambdaEvent is the request API Gateway (a REST API's payload format 1.0, an HTTP API's 1.0 or 2.0) or a function URL (2.0) hands a function: the fields an http.Handler needs of either format.

type LambdaRequestContext

type LambdaRequestContext struct {
	DomainName string `json:"domainName,omitempty"`
	// 1.0's method.
	HTTPMethod string `json:"httpMethod,omitempty"`
	// 2.0's method and path.
	HTTP struct {
		Method string `json:"method,omitempty"`
		Path   string `json:"path,omitempty"`
	} `json:"http"`
}

LambdaRequestContext is what a LambdaEvent says of where it came from.

type LambdaResult

type LambdaResult struct {
	StatusCode        int                 `json:"statusCode"`
	Headers           map[string]string   `json:"headers,omitempty"`
	MultiValueHeaders map[string][]string `json:"multiValueHeaders,omitempty"`
	Cookies           []string            `json:"cookies,omitempty"`
	Body              string              `json:"body"`
	IsBase64Encoded   bool                `json:"isBase64Encoded"`
}

LambdaResult is the proxy result a function answers with. Only the keys the format knows are written, since a REST API refuses a result with others: a header given more than once is in MultiValueHeaders for 1.0, and Set-Cookie in Cookies for 2.0.

type MemoryStore

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

MemoryStore keeps everything in process memory (stores/memory.ts). The default when no store is given, good for tests and for trying the library out. State is gone on restart, so a missed run cannot be noticed across one. Safe for use by many goroutines at once.

func NewMemoryStore

func NewMemoryStore() *MemoryStore

NewMemoryStore is an empty in-memory store.

func (*MemoryStore) Close

func (m *MemoryStore) Close() error

Close does nothing.

func (*MemoryStore) CompareAndSetState

func (m *MemoryStore) CompareAndSetState(_ context.Context, state JobState, expected int64) (bool, error)

func (*MemoryStore) DeleteJob

func (m *MemoryStore) DeleteJob(_ context.Context, name string) error

func (*MemoryStore) DeleteRunIf

func (m *MemoryStore) DeleteRunIf(_ context.Context, id, job string, status RunStatus) (bool, error)

func (*MemoryStore) GetJob

func (m *MemoryStore) GetJob(_ context.Context, name string) (*StoredJob, error)

func (*MemoryStore) GetRun

func (m *MemoryStore) GetRun(_ context.Context, id string) (*Run, error)

func (*MemoryStore) GetState

func (m *MemoryStore) GetState(_ context.Context, job string) (*JobState, error)

func (*MemoryStore) Init

func (m *MemoryStore) Init(context.Context) error

Init does nothing.

func (*MemoryStore) InsertRun

func (m *MemoryStore) InsertRun(_ context.Context, run Run) error

InsertRun refuses an id already recorded, like SQL's primary key.

func (*MemoryStore) LastRun

func (m *MemoryStore) LastRun(ctx context.Context, job string) (*Run, error)

func (*MemoryStore) ListJobs

func (m *MemoryStore) ListJobs(context.Context) ([]StoredJob, error)

ListJobs is every job by name in byte order, as the SQL stores sort.

func (*MemoryStore) ListRuns

func (m *MemoryStore) ListRuns(_ context.Context, job string, limit int) ([]Run, error)

func (*MemoryStore) Prune

func (m *MemoryStore) Prune(_ context.Context, before int64) (int, error)

Prune keeps each job's newest run whatever its age: without it, a job that runs less often than the retention looks like it never ran.

func (*MemoryStore) RunningRuns

func (m *MemoryStore) RunningRuns(context.Context) ([]Run, error)

func (*MemoryStore) SetState

func (m *MemoryStore) SetState(_ context.Context, state JobState) error

func (*MemoryStore) UpdateRun

func (m *MemoryStore) UpdateRun(_ context.Context, run Run) error

UpdateRun changes only the finish's fields; a run that is gone (its job was forgotten) stays gone, as SQL's UPDATE has it.

func (*MemoryStore) UpdateRunIf

func (m *MemoryStore) UpdateRunIf(_ context.Context, run Run, from []RunStatus) (bool, error)

func (*MemoryStore) UpsertJob

func (m *MemoryStore) UpsertJob(_ context.Context, def Definition, now int64) error

type Metric

type Metric struct {
	Name  string
	Value float64
}

Metric is one number a run reported.

type Metrics

type Metrics []Metric

Metrics are a run's numbers in JavaScript's key order: names that are array indices ("10", "200") first in ascending order, then the rest in the order they were first reported. Budgets use the same type.

func (Metrics) Clone

func (m Metrics) Clone() Metrics

Clone is a copy that shares nothing.

func (Metrics) Get

func (m Metrics) Get(name string) (float64, bool)

Get is the value reported for name.

func (Metrics) JSValue

func (m Metrics) JSValue() any

JSValue is the metrics as a JSON object.

func (Metrics) MarshalJSON

func (m Metrics) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

func (*Metrics) Set

func (m *Metrics) Set(name string, value float64)

Set reports a value for name: a later value replaces an earlier one in its place, and a new name takes its place in JavaScript's order.

func (*Metrics) UnmarshalJSON

func (m *Metrics) UnmarshalJSON(b []byte) error

UnmarshalJSON reads a JSON object of numbers.

type MissedDetails

type MissedDetails struct {
	DueAt     int64
	Deadline  float64
	GraceMs   float64
	LastRunAt *int64
}

MissedDetails says which run was missed.

type OpenCondition

type OpenCondition struct {
	Condition Condition
	Since     int64
}

OpenCondition is a condition that is open, and when it opened.

type Option

type Option func(*Client) error

Option configures a client.

func WithAlerts

func WithAlerts(channels ...Channel) Option

WithAlerts is where alerts go, replacing the default console channel. WithAlerts() with none sends nowhere.

func WithClock

func WithClock(now func() int64) Option

WithClock replaces the clock, in epoch milliseconds. Tests use it.

func WithCronSecret

func WithCronSecret(secret string) Option

WithCronSecret is the shared secret job handlers' requests must carry. The default is $CRON_SECRET; "" counts as unset.

func WithDefaults

func WithDefaults(options ...JobOption) Option

WithDefaults applies Grace, Timeout, Timezone and FailuresBeforeAlert to every job that does not set its own.

func WithDeliver

func WithDeliver(mode DeliverMode) Option

WithDeliver sets where alerts are sent from: DeliverNow (the default) or DeliverAtCheck.

func WithErrorHandler

func WithErrorHandler(fn func(err error, where string)) Option

WithErrorHandler is called with anything that goes wrong outside a job: the store failing, an alert channel failing, a triage timeout. where says what was being done ("recording nightly", "alert channel slack"). The default writes to standard error.

func WithRedact

func WithRedact(fn func(text string) string) Option

WithRedact replaces the default redaction (see RedactSecrets) of every run's output and error before it is stored, shown or sent anywhere. A redact function that panics is reported to the error handler and the default is used.

func WithRetention

func WithRetention[D Duration](d D) Option

WithRetention is how long finished runs are kept. Default "30d".

func WithSources

func WithSources(sources ...Source) Option

WithSources adds sources of runs this process does not wrap (the pg_cron source). Each is synced at the start of every check; one that fails is reported to the error handler and the check carries on.

func WithStore

func WithStore(store Store) Option

WithStore is where jobs, runs and state live. The default is a MemoryStore, which forgets on restart.

func WithTriage

func WithTriage(fn TriageFunc) Option

WithTriage adds a short diagnosis to every alert except recoveries.

func WithoutCronSecret

func WithoutCronSecret() Option

WithoutCronSecret lets job handlers run without a secret.

func WithoutRedaction

func WithoutRedaction() Option

WithoutRedaction keeps output and errors exactly as logged.

type OverBudgetDetails

type OverBudgetDetails struct {
	Breaches []BudgetBreach
}

OverBudgetDetails lists the metrics over their limits.

type RecordOption

type RecordOption func(*recordConfig)

RecordOption configures RecordRun.

func WithoutEvaluation

func WithoutEvaluation() RecordOption

WithoutEvaluation stores the run without judging it, for history imported on first sight.

type RecoveredDetails

type RecoveredDetails struct {
	After  []Condition
	Reason string
	Since  *int64
}

RecoveredDetails names the conditions that closed. Reason "unscheduled" closes missed alone because the job no longer has a schedule; Since is when missed opened.

type Routes

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

Routes is the dashboard and JSON API, an http.Handler made by Client.Routes.

func (*Routes) ServeHTTP

func (rt *Routes) ServeHTTP(w http.ResponseWriter, r *http.Request)

ServeHTTP answers one request as the SDK's routes answer it. A store failure (or a panic) is reported to the client's error handler as "routes" and answered 500.

func (*Routes) Token

func (rt *Routes) Token() string

Token is the token the dashboard asks for (the generated one in development), or "" when it is open or locked for want of one.

type RoutesOption

type RoutesOption func(*routesConfig)

RoutesOption configures the dashboard (Client.Routes).

func WithBasePath

func WithBasePath(path string) RoutesOption

WithBasePath is where the dashboard is mounted ("" for the root), so its links resolve. Without it the base is found from the request: what http.StripPrefix took off, else the literal part of the ServeMux pattern that matched (Go 1.22's "/cronwatch/" or "/cronwatch/{path...}"), else DefaultBasePath.

func WithOrigin

func WithOrigin(origin string) RoutesOption

WithOrigin is the public origin the dashboard is served from, such as "https://app.example.com", for an app behind a proxy whose requests carry an internal host or scheme. It stands in for the request's own origin in the cross-site check on writes, the sign-in cookie's Secure flag, the Referer the redirect back after a form follows, and the development sign-in line. Anything that is not an http or https URL is an error from Routes. It takes precedence over WithTrustProxy.

func WithToken

func WithToken(token string) RoutesOption

WithToken is the token the dashboard asks for. Send it as `Authorization: Bearer <token>`, or open the dashboard once with ?token=<token> and a cookie is set. The default is CRONWATCH_TOKEN; "" counts as unset. With no token in development (CRONWATCH_ENV, APP_ENV or GO_ENV naming it), the routes make a random one and print a sign-in link to Stdout on their first request; with no token otherwise they answer 503. /api/check also takes the client's cron secret as a bearer, so a platform cron can run checks without the token.

func WithTrustProxy

func WithTrustProxy() RoutesOption

WithTrustProxy takes the public origin from X-Forwarded-Proto and X-Forwarded-Host (the first value of each, the request's own scheme or host for whichever is missing) when a request carries either. Only for an app whose proxy sets or overwrites both headers: a client can send them too.

func WithoutToken

func WithoutToken() RoutesOption

WithoutToken serves the dashboard to anyone, everywhere, for one behind your own auth (the SDK's `token: null`).

type Run

type Run struct {
	ID         string
	Job        string
	Status     RunStatus
	StartedAt  int64
	FinishedAt *int64
	DurationMs *int64
	Error      *string
	// Lines logged, or the string the job returned. Capped at 16 KB.
	Output  *string
	Metrics Metrics
	// What started the run: "run", "handler", "start" or a value of yours.
	Trigger string
}

Run is one execution of a job, as a store keeps it. Times are epoch milliseconds.

func (Run) JSValue

func (r Run) JSValue() any

JSValue is the run as the SDK writes it.

func (Run) MarshalJSON

func (r Run) MarshalJSON() ([]byte, error)

MarshalJSON writes the SDK's JSON.

func (*Run) UnmarshalJSON

func (r *Run) UnmarshalJSON(b []byte) error

UnmarshalJSON reads the SDK's JSON.

type RunDeleter

type RunDeleter interface {
	// DeleteRunIf deletes the run id only when its stored job is job and its
	// status is status (SQL: DELETE ... WHERE id = ? AND job = ? AND status
	// = ?), and says whether it deleted.
	DeleteRunIf(ctx context.Context, id, job string, status RunStatus) (bool, error)
}

RunDeleter is a store that can take back a run it recorded, only while the run is still of one job and in one status. The SDK has no counterpart: it is how an attempt a queue gave back without failing (a River job that snoozed or cancelled itself, an Asynq task revoked) leaves no run behind, neither a failure nor a success (see DiscardWhen), as the PHP port's stores take back a released Laravel job's attempt.

type RunHandle

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

RunHandle is a run recorded by Start or found by Resume, to finish later. Lines and metrics wait in the handle until Flush or Finish merges them onto a fresh read of the stored run. It is safe for use by many goroutines at once; Flush and Finish take their turns.

func (*RunHandle) Active

func (h *RunHandle) Active() bool

Active is false once finished, and from the start for a resumed run that already finished or does not exist.

func (*RunHandle) Fail

func (h *RunHandle) Fail(ctx context.Context, err error) *Run

Fail finishes the run as failed with err, written like an error a Run function returned.

func (*RunHandle) Finish

func (h *RunHandle) Finish(ctx context.Context) *Run

Finish finishes the run successfully (unless an expect rule says otherwise), judges it like any other and sends what that produces. It returns the run as recorded, or nil when nothing was recorded: the run was already finished (here or elsewhere), was not found, or belongs to another job, which is reported to the error handler. When several processes finish one run, only the one whose write lands judges it. A store that fails is reported, nothing is recorded, and the handle stays active so Finish can be called again.

func (*RunHandle) FinishWith

func (h *RunHandle) FinishWith(ctx context.Context, result any) *Run

FinishWith finishes the run with a result, treated like the value a RunValue function returns: a string is the output when nothing was logged (and what an expect rule checks), and an *http.Response of 400 or more fails the run.

func (*RunHandle) Flush

func (h *RunHandle) Flush(ctx context.Context)

Flush appends the lines and metrics added so far to the stored run, which must still be running and belong to this job. A read, change and write of the run's row, written only while it is still running: two processes appending to one run at the same moment can lose one's lines, but a flush never undoes a finish. Problems go to the error handler.

func (*RunHandle) ID

func (h *RunHandle) ID() string

ID is the run's id.

func (*RunHandle) Job

func (h *RunHandle) Job() string

Job is the job's name.

func (*RunHandle) Log

func (h *RunHandle) Log(parts ...any)

Log adds a line of output, kept in the handle until Flush or Finish.

func (*RunHandle) Metric

func (h *RunHandle) Metric(name string, value float64) error

Metric reports a number for this run. A later value for the same name replaces an earlier one.

func (*RunHandle) StartedAt

func (h *RunHandle) StartedAt() (int64, bool)

StartedAt is when the run started, in epoch milliseconds; false when a resumed run could not be read.

type RunOption

type RunOption func(*runConfig)

RunOption configures one run.

func DiscardWhen

func DiscardWhen(discard func(err error) bool) RunOption

DiscardWhen takes a run back rather than judging it when the function returns an error discard answers true for: an attempt a queue gives back without failing, such as a River job that snoozes itself. The run's row is deleted while it is still running (through the store's RunDeleter), no alert is sent, the job's failures in a row are left as they were, and the error is still returned. A row a check already marked stuck is left as it is, and a store that is not a RunDeleter records the run as it ended; both are reported to the error handler as "discarding <job>". The SDK has no counterpart; the PHP port takes back a released Laravel job's attempt the same way. Run and RunValue only.

func WithRunID

func WithRunID(id string) RunOption

WithRunID gives Start your own stable id for the run, such as a queue's job id: 1 to 200 characters, not starting with "pgcron:" (the pg_cron source's). A start with an id already recorded for this job records nothing and returns a handle on that run instead; an id recorded for another job is an error. Start only.

func WithTrigger

func WithTrigger(trigger string) RunOption

WithTrigger names what started the run. The default is "run" for Run and "start" for Start.

type RunStatus

type RunStatus string

RunStatus is where a run stands.

const (
	StatusRunning RunStatus = "running"
	StatusOK      RunStatus = "ok"
	StatusFailed  RunStatus = "failed"
	StatusTimeout RunStatus = "timeout"
)

The statuses a run moves through.

type RunUpdater

type RunUpdater interface {
	// UpdateRunIf writes the run as UpdateRun does, only when its stored
	// status is one of from, in one step (SQL: UPDATE ... WHERE id = ? AND
	// status IN (...)), and says whether it wrote. This is what lets exactly
	// one of several processes finishing the same run evaluate it.
	UpdateRunIf(ctx context.Context, run Run, from []RunStatus) (bool, error)
}

RunUpdater is a store that can finish a run conditionally. Without it the client reads then writes, which is safe only when one process at a time finishes a given run.

type SlowDetails

type SlowDetails struct {
	DurationMs  int64
	ThresholdMs float64
	// "maxDuration", or how the baseline was worked out.
	Basis string
}

SlowDetails says how slow a run was.

type Source

type Source interface {
	Name() string
	// Sync declares the jobs and records their new runs, returning the
	// alerts recording them sent.
	Sync(ctx context.Context, host SourceHost) ([]Alert, error)
}

Source is where runs this process does not wrap come from, such as pg_cron's jobs. Check syncs each one first, so what it records is evaluated in the same check.

type SourceHost

type SourceHost interface {
	Job(name string, options ...JobOption) (*Job, error)
	RecordRun(ctx context.Context, run Run, options ...RecordOption) ([]Alert, error)
	Store() Store
	Now() int64
	ReportError(err error, where string)
}

SourceHost is what a Source may use of the client. A *Client is one.

type StateComparer

type StateComparer interface {
	// CompareAndSetState writes state only when the stored state's version
	// (absent, or no row at all, counts as 0) equals expected, and says
	// whether it wrote. The client reads, works out the next state, and on
	// a refused write reads again.
	CompareAndSetState(ctx context.Context, state JobState, expected int64) (bool, error)
}

StateComparer is a store that can write a job's state conditionally. Without it the client writes unconditionally, which is safe only when one process at a time writes a job's state.

type Stats

type Stats struct {
	Runs   int
	OkRate float64
	P50Ms  *int64
	P95Ms  *int64
}

Stats summarize a job's last twenty runs of any status; the percentiles are over the successful ones among them.

type Store

type Store interface {
	// Init is called once before first use: create tables here.
	Init(ctx context.Context) error
	UpsertJob(ctx context.Context, definition Definition, now int64) error
	// GetJob is nil when the store does not know the job.
	GetJob(ctx context.Context, name string) (*StoredJob, error)
	// ListJobs is every job, by name in byte order.
	ListJobs(ctx context.Context) ([]StoredJob, error)
	// DeleteJob removes a job, its runs and its state.
	DeleteJob(ctx context.Context, name string) error
	// InsertRun refuses an id already stored.
	InsertRun(ctx context.Context, run Run) error
	// UpdateRun writes a run's status, finish, duration, error, output and
	// metrics. A run that is gone stays gone.
	UpdateRun(ctx context.Context, run Run) error
	// GetRun is nil when there is no such run.
	GetRun(ctx context.Context, id string) (*Run, error)
	// ListRuns is a job's newest runs first, at most limit of them.
	ListRuns(ctx context.Context, job string, limit int) ([]Run, error)
	LastRun(ctx context.Context, job string) (*Run, error)
	// RunningRuns is every run still running, oldest first.
	RunningRuns(ctx context.Context) ([]Run, error)
	// GetState is nil when the job has no state yet.
	GetState(ctx context.Context, job string) (*JobState, error)
	// SetState writes a job's state unconditionally. Used only when the
	// store is not a StateComparer.
	SetState(ctx context.Context, state JobState) error
	// Prune deletes finished runs that started before this time, keeping
	// each job's newest run, and returns how many.
	Prune(ctx context.Context, before int64) (int, error)
	Close() error
}

Store is where jobs, runs and state live. MemoryStore is one; the sqlstore package keeps them in the app's own database. A store of your own should pass the storetest package's contract test.

type StoredJob

type StoredJob struct {
	Name       string
	Definition Definition
	CreatedAt  int64
	UpdatedAt  int64
}

StoredJob is a job as a store knows it.

type TriageContext

type TriageContext struct {
	Alert Alert
	// The job's five newest runs.
	RecentRuns []Run
}

TriageContext is what a triage function is given.

type TriageFunc

type TriageFunc func(ctx context.Context, tc TriageContext) (string, error)

TriageFunc diagnoses an alert in a few sentences, or answers "" for no diagnosis. ctx ends when the client stops waiting (25 seconds); pass it to any request made.

Directories

Path Synopsis
Package alerts holds the SDK's alert channels, request for request, on net/http alone: Slack, Discord and a signed Webhook; the email providers Resend, Postmark, Sendgrid, Mailgun and SES (signed with SigV4, no AWS SDK); Twilio for SMS; and the trackers Sentry, Honeybadger, Datadog, Rollbar, Bugsnag and NewRelic.
Package alerts holds the SDK's alert channels, request for request, on net/http alone: Slack, Discord and a signed Webhook; the email providers Resend, Postmark, Sendgrid, Mailgun and SES (signed with SigV4, no AWS SDK); Twilio for SMS; and the trackers Sentry, Honeybadger, Datadog, Rollbar, Bugsnag and NewRelic.
asynq module
Package bridge is what the scheduler integrations share: the robfigcron, gocron, river and asynq modules beside this one, each a module of its own.
Package bridge is what the scheduler integrations share: the robfigcron, gocron, river and asynq modules beside this one, each a module of its own.
gocron module
internal
js
Package js holds what the port needs of JavaScript's own behaviour, so that every value the SDK writes, compares or counts is written, compared and counted the same way here: numbers as Number.prototype.toString prints them, JSON.stringify and JSON.parse (objects keep JavaScript's key order), string lengths and cuts in UTF-16 code units, the characters \s matches, and Date's calendar arithmetic.
Package js holds what the port needs of JavaScript's own behaviour, so that every value the SDK writes, compares or counts is written, compared and counted the same way here: numbers as Number.prototype.toString prints them, JSON.stringify and JSON.parse (objects keep JavaScript's key order), string lengths and cuts in UTF-16 code units, the characters \s matches, and Date's calendar arithmetic.
jsre
Package jsre is a small backtracking regular expression engine with JavaScript's semantics, for the SDK's secret redaction patterns.
Package jsre is a small backtracking regular expression engine with JavaScript's semantics, for the SDK's secret redaction patterns.
output
Package output is the SDK's output.ts and the recorder of job.ts: the output cap, error text, secret redaction, and the lines and metrics a run collects.
Package output is the SDK's output.ts and the recorder of job.ts: the output cap, error text, secret redaction, and the lines and metrics a run collects.
post
Package post is the one POST the alert channels and Claude triage make, as the SDK makes it with fetch (alerts/shared.ts): the URL cleaned and read as fetch reads it, only http and https, headers checked as fetch checks them, one ten second deadline for the whole request, a redirect refused rather than followed, at most 1 MiB of an answer read, and an error that names only the URL's origin, with every secret the caller holds cut out of a quoted answer before it is cut to 200 characters.
Package post is the one POST the alert channels and Claude triage make, as the SDK makes it with fetch (alerts/shared.ts): the URL cleaned and read as fetch reads it, only http and https, headers checked as fetch checks them, one ten second deadline for the whole request, a redirect refused rather than followed, at most 1 MiB of an answer read, and an error that names only the URL's origin, with every secret the caller holds cut out of a quoted answer before it is cut to 200 characters.
schedule
Package schedule is the SDK's duration.ts and schedule.ts: durations ("15m", "1h30m", a number of milliseconds) parsed and written as the SDK does, and schedules ("0 2 * * *", "@hourly", "every 5m") with their fire times, due times, deadlines and what a run covers.
Package schedule is the SDK's duration.ts and schedule.ts: durations ("15m", "1h30m", a number of milliseconds) parsed and written as the SDK does, and schedules ("0 2 * * *", "@hourly", "every 5m") with their fire times, due times, deadlines and what a run covers.
schedule/cron
Package cron is a port of croner 10, the cron library the SDK uses: its reading of an expression (CronPattern, with its checks and its messages word for word) and its walk to the next matching time (CronDate), habits included: a day the month does not have rolls over, a wall-clock time in a spring-forward gap moves forward by the gap, and a time that happens twice is the earlier one.
Package cron is a port of croner 10, the cron library the SDK uses: its reading of an expression (CronPattern, with its checks and its messages word for word) and its walk to the next matching time (CronDate), habits included: a day the month does not have rolls over, a wall-clock time in a spring-forward gap moves forward by the gap, and a time that happens twice is the earlier one.
webserver command
Command webserver serves the dashboard over HTTP for packages/mcp/test/go-web.test.ts, which drives @cronwatch/mcp against it.
Command webserver serves the dashboard over HTTP for packages/mcp/test/go-web.test.ts, which drives @cronwatch/mcp against it.
Package pgcron watches pg_cron jobs, which run inside Postgres where nothing can wrap them (sources/pgcron.ts).
Package pgcron watches pg_cron jobs, which run inside Postgres where nothing can wrap them (sources/pgcron.ts).
river module
robfigcron module
Package sqlstore keeps CronWatch's jobs, runs and state in the app's own database through database/sql: SQLite, Postgres or MySQL (and MariaDB).
Package sqlstore keeps CronWatch's jobs, runs and state in the app's own database through database/sql: SQLite, Postgres or MySQL (and MariaDB).
Package storetest is the test every CronWatch store passes: the memory store, and the sqlstore package on SQLite, Postgres and MySQL.
Package storetest is the test every CronWatch store passes: the memory store, and the sqlstore package on SQLite, Postgres and MySQL.
Package triage holds Claude triage (triage/anthropic.ts), over plain HTTP: the Messages API is one POST, so no Anthropic SDK is needed.
Package triage holds Claude triage (triage/anthropic.ts), over plain HTTP: the Messages API is one POST, so no Anthropic SDK is needed.

Jump to

Keyboard shortcuts

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