From 9d37d5648653eb5acb09a253e69ccfecd2c5e0f5 Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Wed, 30 Sep 2026 22:35:22 +0200 Subject: [PATCH] feat(11.1-04): add the Services section with a verified Queued jobs page - docs/services/jobs.md: declaring, registering and dispatching jobs, the summer_jobs record, progress, cancellation and workers - conga ExampleJob plus dispatch and status regions run by TestDocsDispatch on the package's Postgres harness (DocsApp in export_docs_test.go) - concept map links the queued jobs row; services/jobs is a required page --- cmd/summer/docs_test.go | 1 + docs/services/jobs.md | 194 ++++++++++++++++++++++++++++ docs/setup/coming-from-wintercms.md | 2 +- docs/site.yaml | 2 + modules/conga/example_test.go | 176 +++++++++++++++++++++++++ modules/conga/export_docs_test.go | 18 +++ 6 files changed, 392 insertions(+), 1 deletion(-) create mode 100644 docs/services/jobs.md create mode 100644 modules/conga/example_test.go create mode 100644 modules/conga/export_docs_test.go diff --git a/cmd/summer/docs_test.go b/cmd/summer/docs_test.go index 176d2e5..a861d8f 100644 --- a/cmd/summer/docs_test.go +++ b/cmd/summer/docs_test.go @@ -93,6 +93,7 @@ var requiredPages = []string{ "console/scaffolding", "console/writing-commands", "console/utilities", + "services/jobs", } // TestDocsRequiredPages asserts that every required page is in the loaded diff --git a/docs/services/jobs.md b/docs/services/jobs.md new file mode 100644 index 0000000..9710e93 --- /dev/null +++ b/docs/services/jobs.md @@ -0,0 +1,194 @@ +--- +title: Queued jobs +description: Declare background jobs with conga.Job, dispatch them inside the caller's transaction, track them in summer_jobs and run workers in serve or on their own. +section: services +order: 110 +--- +# Queued jobs + +WinterCMS pushes slow work onto the Laravel queue and tracks long imports with a job manager. SummerCMS does both with [conga](../../modules/conga/README.md): a plugin declares typed job functions, a caller dispatches them inside its own database transaction, and every dispatched job has a `summer_jobs` row that records its status, progress and outcome. The queue itself is River on the application's Postgres database, so there is no Redis or separate queue server to run. + +Plugin code never imports River. It only uses `conga.Job`, `conga.Manager` and the `pact.HasJobs` interface. + +## Declaring jobs + +A job has two parts: an arguments type and a function. The arguments type implements `pact.JobArgs`: its `Kind` method names the job, and the value is stored as JSON in the queue, so give every field a `json` tag. `conga.Job` wraps a function that takes those arguments into a `pact.Job`. + +```go src=modules/conga/example_test.go#ImportPostsArgs +// ImportPostsArgs are the arguments of the acme.blog post import job. +type ImportPostsArgs struct { + File string `json:"file"` +} +``` + +```go src=modules/conga/example_test.go#ImportPostsArgs.Kind +// Kind names the job. It must be unique across the application. +func (ImportPostsArgs) Kind() string { return "acme_blog_import_posts" } +``` + +```go src=modules/conga/example_test.go#ExampleJob +job := conga.Job(func(ctx context.Context, args ImportPostsArgs) error { + _, dispatched := conga.JobID(ctx) + fmt.Println("import", args.File, "with a summer_jobs row:", dispatched) + return nil +}, conga.OnQueue("imports"), conga.MaxAttempts(5), conga.Timeout(10*time.Minute)) + +// A worker calls Work with the decoded arguments; a unit test can too. +if err := job.Work(context.Background(), ImportPostsArgs{File: "posts.csv"}); err != nil { + fmt.Println(err) +} +fmt.Println(ImportPostsArgs{}.Kind()) +// Output: +// import posts.csv with a summer_jobs row: false +// acme_blog_import_posts +``` + +The options set the job's defaults: + +| Option | Sets | Default | +|--------|------|---------| +| `conga.OnQueue` | The queue the job is inserted on. | `default` | +| `conga.MaxAttempts` | How many times the job is tried before its row is marked as an error. | `queue.max_attempts` (3) | +| `conga.Timeout` | The deadline of each attempt. | `queue.job_timeout` (300 seconds) | + +`conga.JobID` returns the `summer_jobs` row of the running job. It reports `false` when the job has no row, as in the example above, where the function is called directly rather than by a worker. + +## Registering jobs + +A plugin returns its jobs from `pact.HasJobs`. Every worker registers the jobs of every active plugin before it starts, so a job must be declared at build time; there is no way to add one while a worker runs (`conga.ErrRegistrationClosed`). A `pact.Job` that was not built by `conga.Job` is refused with `conga.ErrNotCongaJob`. + +A job that reports progress needs the job manager, so the plugin keeps it from `Register`: + +```go src=modules/conga/example_test.go#BlogPlugin.Jobs +// Jobs returns the plugin's background jobs; every worker registers them. +func (p *BlogPlugin) Jobs() []pact.Job { + return []pact.Job{ImportPostsJob(p.jobs)} +} +``` + +`conga.From` returns the application's `conga.Manager`, publishing one on first use. + +## Dispatching jobs + +Dispatch a job with `conga.Manager.Dispatch`, passing the transaction of the write that needs it. The `summer_jobs` row and the queued job are written on that transaction: when it rolls back, neither exists, and when it commits, a worker picks the job up at once. Laravel gives this guarantee only to jobs marked `afterCommit`; here every dispatch has it. + +```go src=modules/conga/example_test.go#dispatch +m, err := conga.From(app) +if err != nil { + return 0, err +} +var id uint +err = db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + // ... write the import's own records on tx ... + id, err = m.Dispatch(ctx, tx, ImportPostsArgs{File: file}, conga.DispatchOpts{ + Label: "Import posts", + Count: 1, + }) + return err +}) +return id, err +``` + +`conga.DispatchOpts` holds the row's `Label` (required), the initial `Count` for the progress bar and JSON `Metadata`, and can override the job's queue and attempt limit or delay the first attempt with `Delay`. The row also records who dispatched the job: the user ID and admin flag of the authenticated principal in the request context. When the handle you pass is not in a transaction, `Dispatch` opens one of its own. + +For fire-and-forget work that needs no row, `conga.Manager.Enqueue` inserts the job alone, inside the caller's transaction when there is one. + +## The summer_jobs record + +The `summer_jobs` row is what the application reads to show progress. `conga.Manager.Get` returns it as a `conga.Record`, and its `conga.Record.Status` holds the WinterCMS job statuses: + +| Status | Meaning | +|--------|---------| +| `conga.StatusInQueue` | Defined for parity with WinterCMS. `Dispatch` never writes it. | +| `conga.StatusInProgress` | Written by `Dispatch` and kept while the queue retries a failed attempt. | +| `conga.StatusComplete` | The job completed its row, including skipped work recorded with `{"skipped": true}` metadata. | +| `conga.StatusError` | The final attempt failed or panicked; the error text is in the metadata key `error`. | +| `conga.StatusStopped` | The job was cancelled from outside or stopped itself. | + +```go src=modules/conga/example_test.go#status +rec, err := m.Get(ctx, id) +if err != nil { + return "", err +} +switch rec.Status { +case conga.StatusComplete: + return fmt.Sprintf("%s: done (%d/%d)", rec.Label, rec.Progress, rec.ProgressMax), nil +case conga.StatusError, conga.StatusStopped: + return fmt.Sprintf("%s: failed or stopped", rec.Label), nil +default: + return fmt.Sprintf("%s: %d/%d", rec.Label, rec.Progress, rec.ProgressMax), nil +} +``` + +## Progress and cancellation + +A long job reports its progress and honours cancellation between items. The import job below sets the total with `conga.Manager.StartJob`, checks `conga.Manager.CheckIfCanceled` before each item, advances with `conga.Manager.UpdateJobState` and finishes with `conga.Manager.CompleteJob`: + +```go src=modules/conga/example_test.go#ImportPostsJob +// ImportPostsJob is the acme.blog import job. It reports progress on its +// summer_jobs row, stops when the row is cancelled and completes the row. +func ImportPostsJob(m *conga.Manager) pact.Job { + return conga.Job(func(ctx context.Context, args ImportPostsArgs) error { + id, _ := conga.JobID(ctx) + rows := []string{args.File} // ... read the rows of args.File ... + if err := m.StartJob(ctx, id, len(rows)); err != nil { + return err + } + for i := range rows { + canceled, err := m.CheckIfCanceled(ctx, id) + if err != nil { + return err + } + if canceled { + return m.StopJob(ctx, id, nil) + } + // ... import rows[i] ... + if err := m.UpdateJobState(ctx, id, i+1, nil); err != nil { + return err + } + } + return m.CompleteJob(ctx, id, map[string]any{"imported": len(rows)}) + }, conga.OnQueue("imports")) +} +``` + +Cancellation has two sides, as in the WinterCMS job manager: + +- `conga.Manager.CancelJob` is the cancel button. It marks the row as cancelled and stopped, then cancels the queued job, so a job that has not started never runs and a running job's context is cancelled. +- `conga.Manager.StopJob` is what the job calls on its own row after `conga.Manager.CheckIfCanceled` reports `true`. It only sets the stopped status. + +The worker applies the outcome rules around each attempt. An error on an attempt before the last leaves the row in progress so the queue can retry it. The final failed attempt, or a panic on it, marks the row as an error. A job that returns `nil` without completing its row leaves the row as it is, so complete it yourself, as the example does. + +## Running workers + +By default, `serve` runs a worker in the same process as the HTTP server, on every known queue. The known queues are `default`, `scheduled`, every queue in `queue.queues` and every queue a registered job names. + +To run jobs in separate processes, set `queue.work_in_serve` to `false` and start one or more workers: + +```sh +./bin/acme queue:work +./bin/acme queue:work --queue imports --queue default +``` + +`queue:work` runs until it receives SIGINT or SIGTERM, then stops within 10 seconds. An unknown queue name is an error that lists the known ones. + +The worker settings live in `config/queue.yaml`: + +```yaml +work_in_serve: false +max_attempts: 3 +job_timeout: 300 +queues: + default: 4 + imports: 1 +``` + +Each entry under `queues` is the number of jobs of that queue a worker runs at once. The worker listens for new jobs on a dedicated Postgres connection opened from `database.dsn`, so a job committed by any process starts without waiting for the poll interval. With PgBouncer, that connection must use session pooling or go straight to Postgres. + +To delete the waiting jobs of one queue, for example after a bad deploy, run `queue:clear`. Running jobs are never touched: + +```sh +./bin/acme queue:clear imports +``` + +The worker also runs the scheduled console commands that plugins declare. See [Task scheduling](../plugins/scheduling.md). diff --git a/docs/setup/coming-from-wintercms.md b/docs/setup/coming-from-wintercms.md index 5537780..3bed1fb 100644 --- a/docs/setup/coming-from-wintercms.md +++ b/docs/setup/coming-from-wintercms.md @@ -30,7 +30,7 @@ Each row names the WinterCMS concept, the SummerCMS identifiers that replace it | `Event::listen` and `Event::fire` | `festival.Bus.Listen` and `festival.Bus.Fire` on `backpack.App.Events`, keyed by the event's Go type | [festival](../../modules/festival/README.md) | | `App::make` and singleton bindings | `backpack.App.Publish` and `backpack.App.Lookup`, keyed by type | [backpack](../../modules/backpack/README.md) | | Artisan commands and `registerConsoleCommand` | `bonfire.Command` values returned from `pact.HasCommands` | [bonfire](../../modules/bonfire/README.md) | -| Queued jobs | `pact.HasJobs` with jobs built by `conga.Job`, dispatched with `conga.Manager.Dispatch` inside the caller's transaction | [conga](../../modules/conga/README.md) | +| Queued jobs | `pact.HasJobs` with jobs built by `conga.Job`, dispatched with `conga.Manager.Dispatch` inside the caller's transaction | [Queued jobs](../services/jobs.md), [conga](../../modules/conga/README.md) | | `registerSchedule` and the scheduler | `pact.HasSchedule` returning `pact.ScheduledCommand` entries with a `pact.Cadence` | [pact](../../modules/pact/README.md) | | Mail templates in `views/mail` | `pact.HasMailTemplates`, sent through `postcard.Mailer` | [postcard](../../modules/postcard/README.md) | | Settings models and `registerSettings` | `pact.HasSettings` returning `pact.SettingsItem` entries | [cabana](../../modules/cabana/README.md) | diff --git a/docs/site.yaml b/docs/site.yaml index 2ef2702..8ace1cf 100644 --- a/docs/site.yaml +++ b/docs/site.yaml @@ -17,6 +17,8 @@ sections: title: Architecture - name: plugins title: Plugins + - name: services + title: Services - name: console title: Console - name: api diff --git a/modules/conga/example_test.go b/modules/conga/example_test.go new file mode 100644 index 0000000..bb3a3f6 --- /dev/null +++ b/modules/conga/example_test.go @@ -0,0 +1,176 @@ +package conga_test + +import ( + "context" + "fmt" + "testing" + "time" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/conga" + "git.golem15.com/golem15/summercms/modules/pact" + "git.golem15.com/golem15/summercms/modules/party" + "gorm.io/gorm" +) + +// ImportPostsArgs are the arguments of the acme.blog post import job. +type ImportPostsArgs struct { + File string `json:"file"` +} + +// Kind names the job. It must be unique across the application. +func (ImportPostsArgs) Kind() string { return "acme_blog_import_posts" } + +func ExampleJob() { + job := conga.Job(func(ctx context.Context, args ImportPostsArgs) error { + _, dispatched := conga.JobID(ctx) + fmt.Println("import", args.File, "with a summer_jobs row:", dispatched) + return nil + }, conga.OnQueue("imports"), conga.MaxAttempts(5), conga.Timeout(10*time.Minute)) + + // A worker calls Work with the decoded arguments; a unit test can too. + if err := job.Work(context.Background(), ImportPostsArgs{File: "posts.csv"}); err != nil { + fmt.Println(err) + } + fmt.Println(ImportPostsArgs{}.Kind()) + // Output: + // import posts.csv with a summer_jobs row: false + // acme_blog_import_posts +} + +// ImportPostsJob is the acme.blog import job. It reports progress on its +// summer_jobs row, stops when the row is cancelled and completes the row. +func ImportPostsJob(m *conga.Manager) pact.Job { + return conga.Job(func(ctx context.Context, args ImportPostsArgs) error { + id, _ := conga.JobID(ctx) + rows := []string{args.File} // ... read the rows of args.File ... + if err := m.StartJob(ctx, id, len(rows)); err != nil { + return err + } + for i := range rows { + canceled, err := m.CheckIfCanceled(ctx, id) + if err != nil { + return err + } + if canceled { + return m.StopJob(ctx, id, nil) + } + // ... import rows[i] ... + if err := m.UpdateJobState(ctx, id, i+1, nil); err != nil { + return err + } + } + return m.CompleteJob(ctx, id, map[string]any{"imported": len(rows)}) + }, conga.OnQueue("imports")) +} + +// BlogPlugin is the acme.blog plugin; only its jobs are shown here. +type BlogPlugin struct { + jobs *conga.Manager +} + +var _ pact.HasJobs = (*BlogPlugin)(nil) + +func (p *BlogPlugin) ID() string { return "acme.blog" } +func (p *BlogPlugin) Requires() []string { return nil } + +// Register keeps the application's job manager for the plugin's jobs. +func (p *BlogPlugin) Register(app *backpack.App) error { + m, err := conga.From(app) + p.jobs = m + return err +} + +func (p *BlogPlugin) Boot(app *backpack.App) error { return nil } + +// Jobs returns the plugin's background jobs; every worker registers them. +func (p *BlogPlugin) Jobs() []pact.Job { + return []pact.Job{ImportPostsJob(p.jobs)} +} + +// startImport dispatches the import job in the transaction of the write +// that needs it and returns the summer_jobs row id. +func startImport(ctx context.Context, app *backpack.App, db *gorm.DB, file string) (uint, error) { + // docs:start dispatch + m, err := conga.From(app) + if err != nil { + return 0, err + } + var id uint + err = db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + // ... write the import's own records on tx ... + id, err = m.Dispatch(ctx, tx, ImportPostsArgs{File: file}, conga.DispatchOpts{ + Label: "Import posts", + Count: 1, + }) + return err + }) + return id, err + // docs:end dispatch +} + +// importStatus reads the progress of a dispatched job from its row. +func importStatus(ctx context.Context, m *conga.Manager, id uint) (string, error) { + // docs:start status + rec, err := m.Get(ctx, id) + if err != nil { + return "", err + } + switch rec.Status { + case conga.StatusComplete: + return fmt.Sprintf("%s: done (%d/%d)", rec.Label, rec.Progress, rec.ProgressMax), nil + case conga.StatusError, conga.StatusStopped: + return fmt.Sprintf("%s: failed or stopped", rec.Label), nil + default: + return fmt.Sprintf("%s: %d/%d", rec.Label, rec.Progress, rec.ProgressMax), nil + } + // docs:end status +} + +// TestDocsDispatch runs the dispatch region of the Jobs page on the +// package's Postgres harness: a worker registers the plugin's job, the job +// is dispatched in a transaction and completes its row. +func TestDocsDispatch(t *testing.T) { + app, db := conga.DocsApp(t) + plugin := &BlogPlugin{} + if err := plugin.Register(app); err != nil { + t.Fatal(err) + } + if n := len(plugin.Jobs()); n != 1 { + t.Fatalf("Jobs() = %d jobs, want 1", n) + } + w, err := conga.StartWorker(t.Context(), app, []party.Plugin{plugin}, conga.WorkerOptions{}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := w.Stop(ctx); err != nil { + t.Errorf("stop worker: %v", err) + } + }) + + id, err := startImport(t.Context(), app, db, "posts.csv") + if err != nil { + t.Fatal(err) + } + m, err := conga.From(app) + if err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(15 * time.Second) + for { + got, err := importStatus(t.Context(), m, id) + if err != nil { + t.Fatal(err) + } + if got == "Import posts: done (1/1)" { + return + } + if time.Now().After(deadline) { + t.Fatalf("status = %q, want the job to complete", got) + } + time.Sleep(25 * time.Millisecond) + } +} diff --git a/modules/conga/export_docs_test.go b/modules/conga/export_docs_test.go new file mode 100644 index 0000000..37d2dc9 --- /dev/null +++ b/modules/conga/export_docs_test.go @@ -0,0 +1,18 @@ +package conga + +import ( + "testing" + + "git.golem15.com/golem15/summercms/modules/backpack" + "gorm.io/gorm" +) + +// DocsApp returns an app and its GORM handle on a fresh migrated database +// of this package's Postgres harness, for the database-backed docs examples +// in example_test.go. It skips under -short and fails when the harness has +// no database, like the package's other database tests. +func DocsApp(t *testing.T) (*backpack.App, *gorm.DB) { + t.Helper() + db, dsn := migratedDB(t) + return testApp(t, db, dsn, nil) +}