Files
summercms/modules/conga
Jakub Zych 0bc5c77097 feat(11-01): add the conga job framework on River with transactional dispatch
- Pin River v0.47.0 (riverdatabasesql, rivertype) and tidy the example modules
- lagoon.Migrate runs the summercms.conga set: River schema v7 and summer_jobs
  with the apparatus columns plus an internal river_job_id link
- conga.Manager.Dispatch writes the summer_jobs row (status IN_PROGRESS) and the
  River job on the caller's *sql.Tx; a rollback leaves neither
- conga.Job wraps typed job functions so plugins never import River
- conga.StartWorker runs one client on riverdatabasesql.NewWithPgxListener with
  a single-connection LISTEN pool; the final failed attempt sets ERROR
- TestListenPickupLatency: 30s poll, pickup under 1s; poll-only control 2s miss
2026-09-29 15:02:04 +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.

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.Manager.CompleteJob completes a row, conga.Manager.Get reads one, and conga.JobID gives a running job its own row id.
  • 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. 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.

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 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.CompleteJob Sets conga.StatusComplete and progress to progress_max; replaces metadata when given.
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.
conga.WorkerOptions Selects the queues a worker runs.
conga.Worker A running worker; conga.Worker.Stop stops it and conga.Worker.Queues lists its queues.
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.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 Concurrent workers per queue. The worker runs these queues plus every queue a registered job names plus default.
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.
max_attempts: 3
job_timeout: 300
queues:
  default: 4
  imports: 1

Dependencies

  • SummerCMS modules: backpack, 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. go test -short ./modules/conga/... skips the database tests.