# 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:

```go
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:

```go
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:

```go
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:

```go
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:

```go
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`. |

```yaml
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](/docs/api/backpack.md), [bonfire](/docs/api/bonfire.md) (commands and the `bonfire.Catalog` scheduled runs call), [bouncer](/docs/api/bouncer.md) (the dispatching principal), [lagoon](/docs/api/lagoon.md) (the shared pool, `lagoon.JobsTable` and the migrations that create River's schema and `summer_jobs`), [pact](/docs/api/pact.md), [party](/docs/api/party.md).
- 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

```sh
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.
