Files
summercms/modules/conga
Jakub Zych 2237a640d2 feat(11-02): add schedule:run as a scheduler process and a cron --once mode
- schedule:run runs a worker on the scheduled queue with every plugin's periodic jobs
- schedule:run --once runs entries due in the current app.timezone minute without River,
  warns and skips unregistered commands, and returns the first command error
- summer schedule:run delegate forwards --once
- tests for --once minute matching, forged-entry skipping (T-11-09) and ByPeriod dedupe
2026-09-29 22:34:21 +02:00
..

conga

Background jobs on River over the shared Postgres pool: transactional dispatch, a summer_jobs progress record, in-process or dedicated workers, and a wall-clock scheduler.

import "git.golem15.com/golem15/summercms/modules/conga"

Overview

conga is the SummerCMS counterpart of the WinterCMS queue plus the apparatus job manager. Plugins describe background work as typed job functions wrapped by conga.Job and return them from pact.HasJobs, so plugin code never imports River. A caller dispatches a job with conga.Manager.Dispatch inside its own write transaction: the summer_jobs record row and the River job are written on the same *sql.Tx, so a rollback leaves neither behind. River only executes the work; the record row is what progress, outcome and cancellation reads use.

Every River client runs on the one *sql.DB pool that lagoon opens. A worker started by conga.StartWorker uses riverdatabasesql.NewWithPgxListener: all queries go through the shared pool, and only Postgres LISTEN goes through a dedicated pgx pool with a single connection, so a job committed by any process wakes the worker immediately instead of waiting for the poll interval.

The scheduler is the Go form of WinterCMS registerSchedule. Plugins declare recurring console commands through pact.HasSchedule; every worker turns them into River periodic jobs on wall-clock conga.Daily and conga.Every schedules in the app.timezone location. Each due run is a job on the scheduled queue that calls the command in-process through the app's bonfire.Catalog.

Features

  • River-free job declarations: conga.Job turns func(ctx context.Context, args T) error into a pact.Job; conga.OnQueue, conga.MaxAttempts and conga.Timeout set per-job defaults. A pact.Job not built by conga.Job is rejected with conga.ErrNotCongaJob.
  • Transactional dispatch: conga.Manager.Dispatch inserts the summer_jobs row with conga.StatusInProgress, the principal's user id and admin flag, progress_max from conga.DispatchOpts.Count and JSON metadata, then enqueues the River job in the same transaction. It opens a transaction itself when the caller has none.
  • Plain enqueue: conga.Manager.Enqueue inserts a River job without a record row, inside the caller's transaction when there is one.
  • The record row: conga.Record maps summer_jobs; conga.Status holds the WinterCMS status values (conga.StatusInQueue, conga.StatusInProgress, conga.StatusComplete, conga.StatusError, conga.StatusStopped). conga.JobID gives a running job its own row id.
  • The WinterCMS job manager operations with the same semantics: conga.Manager.StartJob, conga.Manager.UpdateJobState, conga.Manager.UpdateMetadata, conga.Manager.CompleteJob, conga.Manager.FailJob, conga.Manager.CheckIfCanceled and conga.Manager.GetMetadata. Updates are raw column writes, so updated_at changes only on dispatch and conga.Manager.StartJob.
  • Cancellation in two parts: conga.Manager.CancelJob is the outside cancel (sets is_canceled and conga.StatusStopped, then cancels the River job, so a queued job never starts and a running job's context is cancelled); conga.Manager.StopJob is what a job calls on its own row after conga.Manager.CheckIfCanceled reports true (status only, the WinterCMS cancelJob).
  • Outcome rules in the worker: an error on an attempt before the last leaves the row in progress so River can retry; the final failed attempt, or a recovered panic on it, sets conga.StatusError with the error text under the metadata key error; an error on a row that was stopped or cancelled cancels the River job instead. A job that returns nil without completing its row leaves it as it is. Skipped work is recorded as complete with {"skipped": true} metadata.
  • Workers: conga.StartWorker registers every plugin job and starts one River client; conga.WorkerOptions.Queues limits it to some queues, and an unknown queue is conga.ErrUnknownQueue listing the known ones. conga.Worker.Stop stops it gracefully and cancels running jobs when its context ends. conga.StartServeWorker is the variant the serve command uses: it starts nothing when queue.work_in_serve is false. An app without jobs still gets a worker that starts and idles.
  • Scheduled commands: every worker carries one River periodic job per pact.HasSchedule entry, in plugin activation order then declaration order, with the id <plugin id>[<index>]:<command>. Only the elected leader enqueues. Each run is a conga.ScheduledCommandArgs job on conga.QueueScheduled with one attempt (an interrupted run is not retried; the next period runs normally) and unique by args within its cadence period, so a leader failover cannot double-enqueue a period. The worker runs a job only when its entry exists in the compiled schedule and its command and arguments match that entry exactly, so a forged river_job row cannot run an arbitrary command. An unregistered command, or an app with no published bonfire.Catalog, is logged at Warn (schedule: command not registered; skipping) and skipped without failing the worker or other entries. Command output is logged line by line at Info with a command attribute; failures are logged with the duration.
  • Wall-clock schedules: conga.Daily fires at the next hh:mm in its location and conga.Every at the next multiple of its interval since local midnight, so a restart never delays a daily run by up to a day the way river.PeriodicInterval(24h) would. On a DST day conga.Daily keeps the wall-clock time. A schedule entry with an empty command, a zero cadence, a daily time out of range, or an interval under one second or not dividing 24h fails the worker start with an error naming the plugin id and entry index.
  • Commands: conga.RuntimeCommands adds queue:work, schedule:run and queue:clear to the application binary. schedule:run runs a scheduler-only worker in its own process, and schedule:run --once runs the entries due in the current minute without River, for system cron.

Usage

A plugin declares a job:

package blog

import (
	"context"

	"git.golem15.com/golem15/summercms/modules/conga"
	"git.golem15.com/golem15/summercms/modules/pact"
)

type ImportPostsArgs struct {
	File string `json:"file"`
}

func (ImportPostsArgs) Kind() string { return "acme_blog_import_posts" }

func ImportPostsJob(m *conga.Manager) pact.Job {
	return conga.Job(func(ctx context.Context, args ImportPostsArgs) error {
		id, _ := conga.JobID(ctx)
		// ... import args.File ...
		return m.CompleteJob(ctx, id, nil)
	}, conga.OnQueue("imports"))
}

and dispatches it inside the write that needs it:

func startImport(ctx context.Context, app *backpack.App, gdb *gorm.DB, file string) (uint, error) {
	m, err := conga.From(app)
	if err != nil {
		return 0, err
	}
	var id uint
	err = gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
		// ... write the import record ...
		id, err = m.Dispatch(ctx, tx, ImportPostsArgs{File: file}, conga.DispatchOpts{Label: "Import posts", Count: 100})
		return err
	})
	return id, err
}

A long job reports progress and honours cancellation between items:

func importPosts(ctx context.Context, m *conga.Manager, files []string) error {
	id, _ := conga.JobID(ctx)
	if err := m.StartJob(ctx, id, len(files)); err != nil {
		return err
	}
	for i, f := range files {
		if canceled, err := m.CheckIfCanceled(ctx, id); err != nil || canceled {
			if err != nil {
				return err
			}
			return m.StopJob(ctx, id, nil)
		}
		// ... import f ...
		_ = f
		if err := m.UpdateJobState(ctx, id, i+1, nil); err != nil {
			return err
		}
	}
	return m.CompleteJob(ctx, id, map[string]any{"imported": len(files)})
}

A plugin schedules one of its registered commands; the worker runs it:

var _ pact.HasSchedule = (*Plugin)(nil)

func (p *Plugin) Schedule() []pact.ScheduledCommand {
	return []pact.ScheduledCommand{
		{Command: "blog:prune-drafts", Cadence: pact.Daily()},
		{Command: "blog:sync-feed", Args: []string{"--quiet"}, Cadence: pact.Every(15 * time.Minute)},
	}
}

A worker runs in the same process or in a separate one:

w, err := conga.StartWorker(ctx, app, plugins, conga.WorkerOptions{})
if err != nil {
	return err
}
defer w.Stop(context.Background())

API reference

Identifier Description
conga.From Returns the app's conga.Manager, publishing one on first use.
conga.Manager App-scoped job manager: registration, dispatch and the summer_jobs record.
conga.Manager.Register Registers jobs built by conga.Job; closed while a worker runs (conga.ErrRegistrationClosed).
conga.Manager.Dispatch Writes the record row and enqueues the River job in one transaction; returns the row id.
conga.Manager.Enqueue Enqueues a River job without a record row.
conga.Manager.StartJob Sets progress to 0, progress_max to the total and updated_at to now.
conga.Manager.UpdateJobState Sets progress; replaces metadata when given.
conga.Manager.UpdateMetadata Replaces metadata.
conga.Manager.CompleteJob Sets conga.StatusComplete and progress to progress_max; replaces metadata when given. Skipped work passes {"skipped": true}.
conga.Manager.FailJob Sets conga.StatusError; replaces metadata when given.
conga.Manager.CancelJob Sets is_canceled and conga.StatusStopped and cancels the River job.
conga.Manager.StopJob Sets conga.StatusStopped only; called by a job on its own row.
conga.Manager.CheckIfCanceled Reports is_canceled.
conga.Manager.GetMetadata Decodes metadata; an empty or non-object value is an empty map.
conga.Manager.Get Reads one conga.Record.
conga.DispatchOpts Label, count, metadata, queue, delay and attempt limit of a dispatch.
conga.EnqueueOpts Queue, delay and attempt limit of an enqueue.
conga.Record The summer_jobs row model.
conga.Status Record status values, matching the WinterCMS job statuses.
conga.Job Wraps a typed job function as a pact.Job that conga can run on River.
conga.JobOption Per-job option: conga.OnQueue, conga.MaxAttempts, conga.Timeout.
conga.JobID Returns the record row id of the job running in a context.
conga.StartWorker Registers plugin jobs and starts a River worker client carrying the plugins' periodic schedule jobs.
conga.StartServeWorker The worker of the serve command; nil when queue.work_in_serve is false.
conga.WorkerOptions Selects the queues a worker runs.
conga.RuntimeCommands Returns the queue:work, schedule:run and queue:clear commands.
conga.Worker A running worker; conga.Worker.Stop stops it and conga.Worker.Queues lists its queues.
conga.Daily river.PeriodicSchedule firing at Hour:Minute every day in Loc (UTC when nil).
conga.Every river.PeriodicSchedule firing at every multiple of Interval since local midnight in Loc (UTC when nil).
conga.ScheduledCommandArgs Args of one scheduled run: compiled Entry id, Command and Args; kind summer.scheduled_command.
conga.QueueScheduled The scheduled queue that scheduled runs are inserted on.
conga.ErrNoDatabase The app has not published the shared database handles.
conga.ErrNotCongaJob A registered pact.Job was not built by conga.Job.
conga.ErrRegistrationClosed Registration was attempted while a worker runs.
conga.ErrUnknownQueue A worker was asked for a queue nothing names.

Configuration

Keys are read from the compass config (config/queue.yaml, or SUMMER_QUEUE__... environment variables).

Key Default Controls
queue.work_in_serve true Whether the serve command runs the job worker in its own process. Set false when a separate queue:work process runs the jobs.
queue.max_attempts 3 Attempts per job before its record becomes conga.StatusError, unless the job or dispatch sets its own.
queue.job_timeout 300 Per-attempt deadline, in seconds or as a duration string such as 5m, unless the job sets conga.Timeout.
queue.queues.<name> default: 4, scheduled: 1 Concurrent workers per queue. The worker runs these queues plus every queue a registered job names plus default and scheduled.
app.timezone UTC IANA location of pact.Daily, pact.DailyAt and pact.Every schedules. An unknown name fails the worker start.
database.dsn none (required) Also opens the worker's single-connection LISTEN pool. With PgBouncer, that connection must use session pooling or go straight to Postgres; transaction pooling cannot hold a LISTEN.
work_in_serve: true
max_attempts: 3
job_timeout: 300
queues:
  default: 4
  imports: 1

CLI commands

conga.RuntimeCommands adds these commands to the application binary. They open the database through lagoon.OpenFromApp, so they need database.dsn and app.key; schedule:run --once opens it only when an entry is due.

Command Arguments and flags Description
queue:work --queue <name>, repeatable Runs a job worker in the foreground on the named queues (default: every known queue) until SIGINT or SIGTERM, then stops it within 10 seconds. An unknown queue is an error that lists the known ones.
schedule:run --once (bare) Without --once: runs a worker on the scheduled queue that carries every plugin's periodic jobs, prints scheduler started and runs until SIGINT or SIGTERM, then stops within 10 seconds. With --once: runs, in-process and without River, every entry due in the current minute of app.timezone (a daily entry at its hour and minute; Every(d) of a minute or less on every run, a longer one when the minutes since midnight are a multiple of d), printing Running scheduled command: <command> <args> per entry or No scheduled commands are ready to run.. An unregistered command prints a warning and is skipped. The exit status is the first command error, after every due entry ran. There is no overlap lock: two runs in one minute run the due entries twice, as Laravel does. System cron: * * * * * cd /app && ./bin/app schedule:run --once.
queue:clear [queue] (default default) Deletes the available, scheduled and retryable jobs of one queue in batches until none are left and prints Cleared N jobs. Running jobs are never touched.

Dependencies

  • SummerCMS modules: backpack, bonfire (commands and the bonfire.Catalog scheduled runs call), bouncer (the dispatching principal), lagoon (the shared pool, lagoon.JobsTable and the migrations that create River's schema and summer_jobs), pact, party.
  • Third-party: github.com/riverqueue/river v0.47.0 with its riverdriver/riverdatabasesql and rivertype modules. River is the Postgres job queue the project stack names: it gives transactional inserts on the shared *sql.DB, retries, stuck-job rescue, leader election for periodic work and LISTEN/NOTIFY wake-ups. github.com/jackc/pgx/v5/pgxpool opens the listener pool; gorm.io/gorm writes the record rows.

Testing

go test ./modules/conga/...

The tests start a postgres:16-alpine container through testcontainers-go and migrate a fresh database per test with lagoon.Migrate, so they need a running Docker daemon. TestListenPickupLatency sets a 30-second poll interval and checks that a job committed by a separate client is picked up in under one second, while a poll-only worker does not pick it up within two. TestScheduleRunsCommand checks that an Every(1s) entry runs its command through a periodic job and that an unregistered command is skipped with a Warn log. go test -short ./modules/conga/... skips the database tests.