From b319e7cc61bdb2c4166b3e4864a1be3ab5f007e4 Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Tue, 29 Sep 2026 15:20:14 +0200 Subject: [PATCH] feat(11-01): run job workers in serve and queue:work, add queue:clear - Manager gains the apparatus JobManager surface: StartJob, UpdateJobState, UpdateMetadata, FailJob, CancelJob (is_canceled + STOPPED + River JobCancel), StopJob (STOPPED only), CheckIfCanceled and GetMetadata, all raw column writes so updated_at is untouched - serve starts the in-process worker unless queue.work_in_serve is false and stops it on shutdown; an app without jobs gets an idle worker - queue:work runs a foreground worker with repeatable --queue filters; queue:clear deletes available, scheduled and retryable jobs of one queue - the generated main appends conga.RuntimeCommands; summer delegates queue:work and queue:clear; make:job scaffolds a conga.Job --- cmd/summer/main.go | 2 + cmd/summer/main_test.go | 2 +- cmd/summer/runtime.go | 20 ++ examples/hello/main.go | 5 + internal/build/artifact.go | 10 +- internal/build/build.go | 2 + internal/build/build_test.go | 20 +- internal/build/stubs/artifacts.tmpl | 17 +- modules/conga/README.md | 57 ++- modules/conga/commands.go | 118 +++++++ modules/conga/conga.go | 119 +++++++ modules/conga/listen_test.go | 519 ++++++++++++++++++++++++++++ modules/conga/worker.go | 34 ++ modules/surf/README.md | 6 +- modules/surf/serve.go | 46 ++- 15 files changed, 939 insertions(+), 38 deletions(-) create mode 100644 modules/conga/commands.go diff --git a/cmd/summer/main.go b/cmd/summer/main.go index e0fcc89..2cc8b2c 100644 --- a/cmd/summer/main.go +++ b/cmd/summer/main.go @@ -42,6 +42,8 @@ func toolCommands() []bonfire.Command { delegateRollbackCommand(), delegateCommand("migrate:status", "Show per-plugin migration history in the app binary"), delegateCommand("serve", "Run the app HTTP server"), + delegateQueueWorkCommand(), + delegateCommand("queue:clear", "Clear pending queued jobs in the app binary"), } } diff --git a/cmd/summer/main_test.go b/cmd/summer/main_test.go index e7a8711..34fb858 100644 --- a/cmd/summer/main_test.go +++ b/cmd/summer/main_test.go @@ -19,7 +19,7 @@ func TestToolCommandNames(t *testing.T) { for _, c := range toolCommands() { names = append(names, c.Name) } - for _, want := range []string{"build", "make:plugin", "make:model", "make:migration", "make:command", "make:job", "make:admin-controller", "plugin:add", "dev", "migrate", "migrate:rollback", "migrate:status", "serve"} { + for _, want := range []string{"build", "make:plugin", "make:model", "make:migration", "make:command", "make:job", "make:admin-controller", "plugin:add", "dev", "migrate", "migrate:rollback", "migrate:status", "serve", "queue:work", "queue:clear"} { if !slices.Contains(names, want) { t.Fatalf("missing %s in %v", want, names) } diff --git a/cmd/summer/runtime.go b/cmd/summer/runtime.go index ae91256..a916dd6 100644 --- a/cmd/summer/runtime.go +++ b/cmd/summer/runtime.go @@ -43,6 +43,26 @@ func delegateRollbackCommand() bonfire.Command { } } +func delegateQueueWorkCommand() bonfire.Command { + return bonfire.Command{ + Name: "queue:work", + Description: "Run background job workers in the app binary", + Flags: []bonfire.Flag{{ + Name: "queue", + Description: "Queue to work (repeatable; default all known queues)", + Repeatable: true, + }}, + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + argv := []string{"queue:work"} + for _, q := range in.Flags("queue") { + argv = append(argv, "--queue", q) + } + argv = append(argv, in.Args()...) + return runAppBinary(ctx, argv, out) + }, + } +} + func runAppBinary(ctx context.Context, argv []string, out bonfire.Output) error { if len(argv) == 0 { return fmt.Errorf("summer: missing app command") diff --git a/examples/hello/main.go b/examples/hello/main.go index 8968153..3eceaf7 100644 --- a/examples/hello/main.go +++ b/examples/hello/main.go @@ -9,7 +9,9 @@ import ( "git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/bonfire" + "git.golem15.com/golem15/summercms/modules/cabana" "git.golem15.com/golem15/summercms/modules/compass" + "git.golem15.com/golem15/summercms/modules/conga" "git.golem15.com/golem15/summercms/modules/lagoon" "git.golem15.com/golem15/summercms/modules/pact" "git.golem15.com/golem15/summercms/modules/party" @@ -34,7 +36,10 @@ func run(args []string, out io.Writer) error { return err } commands := lagoon.RuntimeCommands(app, plugins) + commands = append(commands, conga.RuntimeCommands(app, plugins)...) commands = append(commands, surf.ServeCommand(app, plugins)) + commands = append(commands, surf.RouteListCommand(app, plugins)) + commands = append(commands, cabana.RuntimeCommands(app)...) for _, plugin := range plugins { if hasCommands, ok := plugin.(pact.HasCommands); ok { commands = append(commands, hasCommands.Commands()...) diff --git a/internal/build/artifact.go b/internal/build/artifact.go index 5041625..d78596a 100644 --- a/internal/build/artifact.go +++ b/internal/build/artifact.go @@ -192,7 +192,8 @@ func MakeCommand(ctx context.Context, startDir, pluginID, name string) (Artifact return ArtifactResult{PluginDir: plugin.Dir, Files: []string{path}, Hint: plugin.Hint}, nil } -// MakeJob writes a pact.Job stub without importing River. +// MakeJob writes a typed job wrapped in conga.Job; the plugin never +// imports River. func MakeJob(ctx context.Context, startDir, pluginID, name string) (ArtifactResult, error) { plugin, err := resolvePlugin(startDir, pluginID) if err != nil { @@ -210,10 +211,9 @@ func MakeJob(ctx context.Context, startDir, pluginID, name string) (ArtifactResu return ArtifactResult{}, err } src, err := renderStub("job.go", artifactData{ - Ident: name, - Func: fn, - Worker: unexportedIdent(name) + "Job", - Kind: plugin.ID + "." + toSnake(name), + Ident: name, + Func: fn, + Kind: plugin.ID + "." + toSnake(name), }) if err != nil { return ArtifactResult{}, err diff --git a/internal/build/build.go b/internal/build/build.go index dc5f329..ca57540 100644 --- a/internal/build/build.go +++ b/internal/build/build.go @@ -89,6 +89,7 @@ func generateMain(m Manifest) ([]byte, error) { b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/bonfire\"\n") b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/cabana\"\n") b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/compass\"\n") + b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/conga\"\n") b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/lagoon\"\n") b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/pact\"\n") b.WriteString("\t\"git.golem15.com/golem15/summercms/modules/party\"\n") @@ -111,6 +112,7 @@ func generateMain(m Manifest) ([]byte, error) { b.WriteString("\t\treturn err\n") b.WriteString("\t}\n") b.WriteString("\tcommands := lagoon.RuntimeCommands(app, plugins)\n") + b.WriteString("\tcommands = append(commands, conga.RuntimeCommands(app, plugins)...)\n") b.WriteString("\tcommands = append(commands, surf.ServeCommand(app, plugins))\n") b.WriteString("\tcommands = append(commands, surf.RouteListCommand(app, plugins))\n") b.WriteString("\tcommands = append(commands, cabana.RuntimeCommands(app)...)\n") diff --git a/internal/build/build_test.go b/internal/build/build_test.go index 3565250..1b7d36b 100644 --- a/internal/build/build_test.go +++ b/internal/build/build_test.go @@ -128,6 +128,20 @@ func TestGenerateMainRegistersCabanaRuntimeCommands(t *testing.T) { } } +func TestGenerateMainRegistersCongaRuntimeCommands(t *testing.T) { + mainSrc, err := generateMain(Manifest{Module: "example.com/app", Binary: "hello"}) + if err != nil { + t.Fatal(err) + } + line := []byte("commands = append(commands, conga.RuntimeCommands(app, plugins)...)") + if n := bytes.Count(mainSrc, line); n != 1 { + t.Fatalf("conga.RuntimeCommands line count = %d, want 1\n%s", n, mainSrc) + } + if !bytes.Contains(mainSrc, []byte(`"git.golem15.com/golem15/summercms/modules/conga"`)) { + t.Fatalf("generated main does not import conga:\n%s", mainSrc) + } +} + func TestWriteIfChangedSkipsIdenticalBytes(t *testing.T) { dir := t.TempDir() path := filepath.Join(dir, "out.go") @@ -696,9 +710,9 @@ func TestScaffoldAllArtifacts(t *testing.T) { for _, want := range []string{ "func (ReindexArgs) Kind()", `Kind() string { return "golem15.demo.reindex" }`, - "func (reindexJob) Work(", - "unexpected args type", - "func ReindexJob()", + "func ReindexJob() pact.Job", + "conga.Job(func(ctx context.Context, args ReindexArgs) error", + `"git.golem15.com/golem15/summercms/modules/conga"`, } { if !bytes.Contains(jobSrc, []byte(want)) { t.Fatalf("job.go missing %s:\n%s", want, jobSrc) diff --git a/internal/build/stubs/artifacts.tmpl b/internal/build/stubs/artifacts.tmpl index 950d10f..7934486 100644 --- a/internal/build/stubs/artifacts.tmpl +++ b/internal/build/stubs/artifacts.tmpl @@ -66,8 +66,8 @@ package jobs import ( "context" - "fmt" + "git.golem15.com/golem15/summercms/modules/conga" "git.golem15.com/golem15/summercms/modules/pact" ) @@ -76,17 +76,12 @@ type {{.Ident}}Args struct{} func ({{.Ident}}Args) Kind() string { return {{printf "%q" .Kind}} } -type {{.Worker}} struct{} - -func ({{.Worker}}) Work(ctx context.Context, args pact.JobArgs) error { - if _, ok := args.({{.Ident}}Args); !ok { - return fmt.Errorf("jobs: unexpected args type %T", args) - } - return nil -} - // {{.Func}} returns the {{.Kind}} job. -func {{.Func}}() pact.Job { return {{.Worker}}{} } +func {{.Func}}() pact.Job { + return conga.Job(func(ctx context.Context, args {{.Ident}}Args) error { + return nil + }) +} {{end}} {{define "admin_controller.go"}}// Code generated by summer make. DO NOT EDIT. diff --git a/modules/conga/README.md b/modules/conga/README.md index 701e16d..ef0807c 100644 --- a/modules/conga/README.md +++ b/modules/conga/README.md @@ -15,9 +15,12 @@ Every River client runs on the one `*sql.DB` pool that `lagoon` opens. A worker - 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. +- 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. +- Commands: `conga.RuntimeCommands` adds `queue:work` and `queue:clear` to the application binary. ## Usage @@ -66,6 +69,31 @@ func startImport(ctx context.Context, app *backpack.App, gdb *gorm.DB, file stri } ``` +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 worker runs in the same process or in a separate one: ```go @@ -85,7 +113,15 @@ defer w.Stop(context.Background()) | `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.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. | @@ -95,7 +131,9 @@ defer w.Stop(context.Background()) | `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.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` and `queue:clear` commands. | | `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`. | @@ -108,12 +146,14 @@ Keys are read from the compass config (`config/queue.yaml`, or `SUMMER_QUEUE__.. | 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.` | `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`. | ```yaml +work_in_serve: true max_attempts: 3 job_timeout: 300 queues: @@ -121,6 +161,15 @@ queues: imports: 1 ``` +## CLI commands + +`conga.RuntimeCommands` adds these commands to the application binary. Both open the database through `lagoon.OpenFromApp`, so they need `database.dsn` and `app.key`. + +| Command | Arguments and flags | Description | +|---------|---------------------|-------------| +| `queue:work` | `--queue `, 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. | +| `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](../backpack/README.md), [bouncer](../bouncer/README.md) (the dispatching principal), [lagoon](../lagoon/README.md) (the shared pool, `lagoon.JobsTable` and the migrations that create River's schema and `summer_jobs`), [pact](../pact/README.md), [party](../party/README.md). diff --git a/modules/conga/commands.go b/modules/conga/commands.go new file mode 100644 index 0000000..b352be1 --- /dev/null +++ b/modules/conga/commands.go @@ -0,0 +1,118 @@ +package conga + +import ( + "context" + "fmt" + "os" + "os/signal" + "strings" + "syscall" + "time" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/bonfire" + "git.golem15.com/golem15/summercms/modules/lagoon" + "git.golem15.com/golem15/summercms/modules/party" + "github.com/riverqueue/river" + "github.com/riverqueue/river/rivertype" +) + +// clearBatch is the number of jobs one JobDeleteMany pass removes. +const clearBatch = 10000 + +// RuntimeCommands returns the queue:work and queue:clear commands for the +// application binary. +func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Command { + return []bonfire.Command{ + { + Name: "queue:work", + Description: "Run background job workers in the foreground", + Flags: []bonfire.Flag{{ + Name: "queue", + Description: "Queue to work (repeatable; default all known queues)", + Repeatable: true, + }}, + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + return withDB(ctx, app, func() error { + return work(ctx, app, plugins, in.Flags("queue"), out) + }) + }, + }, + { + Name: "queue:clear", + Description: "Clear all queued jobs, by deleting all pending jobs.", + Args: []bonfire.Arg{{ + Name: "queue", + Description: "The name of the queue to clear (default \"default\")", + }}, + Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + queue, _ := in.Argument("queue") + queue = strings.TrimSpace(queue) + if queue == "" { + queue = defaultQueue + } + return withDB(ctx, app, func() error { + return clearQueue(ctx, app, queue, out) + }) + }, + }, + } +} + +func work(ctx context.Context, app *backpack.App, plugins []party.Plugin, queues []string, out bonfire.Output) error { + ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) + defer stop() + w, err := StartWorker(ctx, app, plugins, WorkerOptions{Queues: queues}) + if err != nil { + return err + } + out.Info("worker started on queues: " + strings.Join(w.Queues(), ", ")) + <-ctx.Done() + stopCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return w.Stop(stopCtx) +} + +// clearQueue deletes the available, scheduled and retryable jobs of queue in +// batches until a pass deletes none. River never deletes running jobs. +func clearQueue(ctx context.Context, app *backpack.App, queue string, out bonfire.Output) error { + m, err := From(app) + if err != nil { + return err + } + client, err := m.insertClient() + if err != nil { + return err + } + out.Info(fmt.Sprintf("Clearing queue %q", queue)) + total := 0 + for { + res, err := client.JobDeleteMany(ctx, river.NewJobDeleteManyParams(). + Queues(queue). + States(rivertype.JobStateAvailable, rivertype.JobStateScheduled, rivertype.JobStateRetryable). + First(clearBatch)) + if err != nil { + return fmt.Errorf("conga: clear queue %q: %w", queue, err) + } + if len(res.Jobs) == 0 { + break + } + total += len(res.Jobs) + } + out.Info(fmt.Sprintf("Cleared %d jobs", total)) + return nil +} + +// withDB opens and publishes the shared pool for a CLI command and closes it +// when fn returns, the way lagoon's own commands do. +func withDB(ctx context.Context, app *backpack.App, fn func() error) error { + sqlDB, gdb, err := lagoon.OpenFromApp(ctx, app) + if err != nil { + return err + } + defer sqlDB.Close() + if err := lagoon.Publish(app, sqlDB, gdb); err != nil { + return err + } + return fn() +} diff --git a/modules/conga/conga.go b/modules/conga/conga.go index cbe889c..c8fd321 100644 --- a/modules/conga/conga.go +++ b/modules/conga/conga.go @@ -18,6 +18,7 @@ import ( "git.golem15.com/golem15/summercms/modules/lagoon" "git.golem15.com/golem15/summercms/modules/pact" "github.com/riverqueue/river" + "github.com/riverqueue/river/rivertype" "gorm.io/gorm" ) @@ -256,6 +257,124 @@ func (m *Manager) CompleteJob(ctx context.Context, id uint, metadata map[string] return m.updateColumns(ctx, gdb, id, cols) } +// StartJob sets progress to 0, progress_max to total and updated_at to now. +func (m *Manager) StartJob(ctx context.Context, id uint, total int) error { + gdb, err := m.gdb() + if err != nil { + return err + } + return m.updateColumns(ctx, gdb, id, map[string]any{"progress": 0, "progress_max": total, "updated_at": time.Now()}) +} + +// UpdateJobState sets progress only, then replaces metadata when it is +// non-empty. updated_at is not touched. +func (m *Manager) UpdateJobState(ctx context.Context, id uint, current int, metadata map[string]any) error { + gdb, err := m.gdb() + if err != nil { + return err + } + if err := m.updateColumns(ctx, gdb, id, map[string]any{"progress": current}); err != nil { + return err + } + if len(metadata) > 0 { + return m.UpdateMetadata(ctx, id, metadata) + } + return nil +} + +// UpdateMetadata replaces the row's metadata. updated_at is not touched. +func (m *Manager) UpdateMetadata(ctx context.Context, id uint, metadata map[string]any) error { + gdb, err := m.gdb() + if err != nil { + return err + } + enc, err := encodeMetadata(metadata) + if err != nil { + return err + } + return m.updateColumns(ctx, gdb, id, map[string]any{"metadata": enc}) +} + +// FailJob sets StatusError, replacing metadata only when it is non-empty. +// Job functions normally just return an error; the worker records +// StatusError on the final attempt. +func (m *Manager) FailJob(ctx context.Context, id uint, metadata map[string]any) error { + return m.setStatus(ctx, id, StatusError, metadata) +} + +// StopJob sets StatusStopped only, the WinterCMS cancelJob semantics. A job +// calls it on its own row after CheckIfCanceled reports true, then returns. +// It neither sets is_canceled nor touches the River job. +func (m *Manager) StopJob(ctx context.Context, id uint, metadata map[string]any) error { + return m.setStatus(ctx, id, StatusStopped, metadata) +} + +// CancelJob cancels a job from outside it: it sets is_canceled and +// StatusStopped in one update, then cancels the River job, so a queued job +// never starts and a running job's ctx is cancelled. A River job that has +// already finished or been removed is not an error. +func (m *Manager) CancelJob(ctx context.Context, id uint) error { + gdb, err := m.gdb() + if err != nil { + return err + } + if err := m.updateColumns(ctx, gdb, id, map[string]any{"is_canceled": true, "status": int(StatusStopped)}); err != nil { + return err + } + rec, err := m.Get(ctx, id) + if err != nil { + return err + } + if rec.RiverJobID == nil { + return nil + } + client, err := m.insertClient() + if err != nil { + return err + } + if _, err := client.JobCancel(ctx, *rec.RiverJobID); err != nil && !errors.Is(err, rivertype.ErrNotFound) { + return fmt.Errorf("conga: cancel river job %d: %w", *rec.RiverJobID, err) + } + return nil +} + +// CheckIfCanceled reports the row's is_canceled flag. Long jobs call it +// between items and stop with StopJob when it is true. +func (m *Manager) CheckIfCanceled(ctx context.Context, id uint) (bool, error) { + rec, err := m.Get(ctx, id) + if err != nil { + return false, err + } + return rec.IsCanceled, nil +} + +// GetMetadata decodes the row's metadata. Anything that is not a non-empty +// JSON object, including the "" written by a dispatch without metadata, +// decodes to an empty map. +func (m *Manager) GetMetadata(ctx context.Context, id uint) (map[string]any, error) { + gdb, err := m.gdb() + if err != nil { + return nil, err + } + return m.metadata(ctx, gdb, id) +} + +func (m *Manager) setStatus(ctx context.Context, id uint, status Status, metadata map[string]any) error { + gdb, err := m.gdb() + if err != nil { + return err + } + cols := map[string]any{"status": int(status)} + if len(metadata) > 0 { + enc, err := encodeMetadata(metadata) + if err != nil { + return err + } + cols["metadata"] = enc + } + return m.updateColumns(ctx, gdb, id, cols) +} + // Get returns the summer_jobs row with id. func (m *Manager) Get(ctx context.Context, id uint) (Record, error) { gdb, err := m.gdb() diff --git a/modules/conga/listen_test.go b/modules/conga/listen_test.go index a8c96c0..6171b44 100644 --- a/modules/conga/listen_test.go +++ b/modules/conga/listen_test.go @@ -1,18 +1,25 @@ package conga import ( + "bytes" "context" "database/sql" "errors" + "reflect" + "strings" + "sync/atomic" "testing" "time" "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/bonfire" "git.golem15.com/golem15/summercms/modules/bouncer" + "git.golem15.com/golem15/summercms/modules/compass" "git.golem15.com/golem15/summercms/modules/pact" "git.golem15.com/golem15/summercms/modules/party" "github.com/riverqueue/river" "github.com/riverqueue/river/riverdriver/riverdatabasesql" + "github.com/riverqueue/river/rivertype" "gorm.io/gorm" ) @@ -267,3 +274,515 @@ func TestDispatchTransactional(t *testing.T) { } }) } + +// testAppKey is a fixed test-only app.key (32 zero bytes). +const testAppKey = "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA=" + +type failArgs struct { + Mode string `json:"mode"` +} + +func (failArgs) Kind() string { return "conga_test_fail" } + +type blockArgs struct{} + +func (blockArgs) Kind() string { return "conga_test_block" } + +// immediateRetry schedules a retry right away so retry tests stay fast. +type immediateRetry struct{} + +func (immediateRetry) NextRetry(*rivertype.JobRow) time.Time { + return time.Now().Add(10 * time.Millisecond) +} + +// commandApp is an app with config only, the way the generated main builds +// it before a CLI command opens the database. +func commandApp(t *testing.T, dsn string, extra map[string]any) *backpack.App { + t.Helper() + cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}}) + if err != nil { + t.Fatal(err) + } + if err := cfg.Set("database.dsn", dsn); err != nil { + t.Fatal(err) + } + if err := cfg.Set("app.key", testAppKey); err != nil { + t.Fatal(err) + } + for k, v := range extra { + if err := cfg.Set(k, v); err != nil { + t.Fatal(err) + } + } + return backpack.New(cfg) +} + +func riverState(t *testing.T, gdb *gorm.DB, riverID *int64) string { + t.Helper() + if riverID == nil { + t.Fatal("record has no river_job_id") + } + var state string + if err := gdb.Raw(`SELECT state::text FROM river_job WHERE id = ?`, *riverID).Scan(&state).Error; err != nil { + t.Fatal(err) + } + return state +} + +func waitRiverState(t *testing.T, gdb *gorm.DB, riverID *int64, want string) { + t.Helper() + deadline := time.Now().Add(10 * time.Second) + got := "" + for time.Now().Before(deadline) { + if got = riverState(t, gdb, riverID); got == want { + return + } + time.Sleep(25 * time.Millisecond) + } + t.Fatalf("river job state = %q, want %q", got, want) +} + +func mustGet(t *testing.T, m *Manager, id uint) Record { + t.Helper() + rec, err := m.Get(t.Context(), id) + if err != nil { + t.Fatal(err) + } + return rec +} + +func mustMetadata(t *testing.T, m *Manager, id uint) map[string]any { + t.Helper() + meta, err := m.GetMetadata(t.Context(), id) + if err != nil { + t.Fatal(err) + } + return meta +} + +// TestJobManagerOperations covers the WinterCMS JobManager methods on the +// summer_jobs row (D-02, D-03, D-04). +func TestJobManagerOperations(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + ctx := t.Context() + dispatch := func(label string) uint { + t.Helper() + id, err := m.Dispatch(ctx, gdb, pingArgs{Tag: label}, DispatchOpts{Label: label}) + if err != nil { + t.Fatal(err) + } + return id + } + old := time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC) + ageRow := func(id uint) { + t.Helper() + if err := gdb.Exec(`UPDATE summer_jobs SET updated_at = ? WHERE id = ?`, old, id).Error; err != nil { + t.Fatal(err) + } + } + untouched := func(id uint) { + t.Helper() + if rec := mustGet(t, m, id); rec.UpdatedAt == nil || !rec.UpdatedAt.Equal(old) { + t.Fatalf("updated_at = %v, want untouched %v", rec.UpdatedAt, old) + } + } + + id := dispatch("Ops") + t.Run("start", func(t *testing.T) { + ageRow(id) + if err := m.StartJob(ctx, id, 5); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + if rec.Progress != 0 || rec.ProgressMax != 5 { + t.Fatalf("progress = %d/%d, want 0/5", rec.Progress, rec.ProgressMax) + } + if rec.UpdatedAt == nil || rec.UpdatedAt.Equal(old) { + t.Fatal("StartJob did not set updated_at") + } + }) + t.Run("update_state", func(t *testing.T) { + ageRow(id) + if err := m.UpdateJobState(ctx, id, 3, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Progress != 3 || rec.Metadata != `""` { + t.Fatalf("progress/metadata = %d/%q, want 3/\"\"", rec.Progress, rec.Metadata) + } + if err := m.UpdateJobState(ctx, id, 4, map[string]any{"a": 1}); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Progress != 4 { + t.Fatalf("progress = %d, want 4", rec.Progress) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"a": float64(1)}) { + t.Fatalf("metadata = %v", got) + } + untouched(id) + }) + t.Run("update_metadata", func(t *testing.T) { + ageRow(id) + if err := m.UpdateMetadata(ctx, id, map[string]any{"b": "x"}); err != nil { + t.Fatal(err) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { + t.Fatalf("metadata = %v", got) + } + untouched(id) + }) + t.Run("get_metadata_of_empty_string", func(t *testing.T) { + got := mustMetadata(t, m, dispatch("Empty")) + if got == nil || len(got) != 0 { + t.Fatalf("metadata = %#v, want empty map", got) + } + }) + t.Run("fail", func(t *testing.T) { + if err := m.FailJob(ctx, id, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Status != StatusError { + t.Fatalf("status = %d, want %d", rec.Status, StatusError) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { + t.Fatalf("FailJob(nil) replaced metadata: %v", got) + } + if err := m.FailJob(ctx, id, map[string]any{"x": 1}); err != nil { + t.Fatal(err) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"x": float64(1)}) { + t.Fatalf("metadata = %v", got) + } + }) + t.Run("stop", func(t *testing.T) { + sid := dispatch("Stop") + if err := m.StopJob(ctx, sid, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, sid); rec.Status != StatusStopped || rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/false", rec.Status, rec.IsCanceled) + } + if canceled, err := m.CheckIfCanceled(ctx, sid); err != nil || canceled { + t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) + } + }) + t.Run("cancel", func(t *testing.T) { + cid := dispatch("Cancel") + if err := m.CancelJob(ctx, cid); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, cid) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) + } + if canceled, err := m.CheckIfCanceled(ctx, cid); err != nil || !canceled { + t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) + } + if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { + t.Fatalf("river state = %q, want cancelled", state) + } + }) + t.Run("skip", func(t *testing.T) { + kid := dispatch("Skip") + if err := m.CompleteJob(ctx, kid, map[string]any{"skipped": true}); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, kid); rec.Status != StatusComplete { + t.Fatalf("status = %d, want %d", rec.Status, StatusComplete) + } + if got := mustMetadata(t, m, kid); !reflect.DeepEqual(got, map[string]any{"skipped": true}) { + t.Fatalf("metadata = %v", got) + } + }) +} + +// TestAttemptOutcomes covers D-03: retries keep the row in progress and only +// the final failed attempt, or a panic on it, records StatusError. +func TestAttemptOutcomes(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + seen := make(chan Status, 16) + job := Job(func(ctx context.Context, a failArgs) error { + if id, ok := JobID(ctx); ok { + if rec, err := m.Get(ctx, id); err == nil { + seen <- rec.Status + } + } + if a.Mode == "panic" { + panic("boom") + } + return errors.New("import failed") + }) + startTestWorker(t, app, WorkerOptions{pollInterval: 100 * time.Millisecond, retryPolicy: immediateRetry{}}, job) + ctx := t.Context() + + t.Run("retries_then_error", func(t *testing.T) { + id, err := m.Dispatch(ctx, gdb, failArgs{Mode: "error"}, DispatchOpts{Label: "Retry", MaxAttempts: 3}) + if err != nil { + t.Fatal(err) + } + rec := waitStatus(t, m, id, StatusError) + for attempt := 1; attempt <= 3; attempt++ { + select { + case st := <-seen: + if st != StatusInProgress { + t.Fatalf("status at the start of attempt %d = %d, want %d", attempt, st, StatusInProgress) + } + case <-time.After(5 * time.Second): + t.Fatalf("attempt %d did not run", attempt) + } + } + if got := mustMetadata(t, m, id); got["error"] != "import failed" { + t.Fatalf("metadata = %v, want error text", got) + } + waitRiverState(t, gdb, rec.RiverJobID, "discarded") + }) + + t.Run("panic", func(t *testing.T) { + id, err := m.Dispatch(ctx, gdb, failArgs{Mode: "panic"}, DispatchOpts{Label: "Panic", MaxAttempts: 1}) + if err != nil { + t.Fatal(err) + } + waitStatus(t, m, id, StatusError) + if got := mustMetadata(t, m, id); !strings.Contains(got["error"].(string), "boom") { + t.Fatalf("metadata = %v, want the panic text", got) + } + }) +} + +// TestCancelJob covers D-04: a queued job never runs and a running job's ctx +// is cancelled while its row stays STOPPED. +func TestCancelJob(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + started := make(chan uint, 4) + var runs atomic.Int32 + job := Job(func(ctx context.Context, a blockArgs) error { + runs.Add(1) + id, _ := JobID(ctx) + started <- id + <-ctx.Done() + return ctx.Err() + }, Timeout(30*time.Second)) + startTestWorker(t, app, WorkerOptions{}, job) + ctx := t.Context() + + t.Run("queued", func(t *testing.T) { + id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Queued", Delay: time.Hour}) + if err != nil { + t.Fatal(err) + } + if err := m.CancelJob(ctx, id); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) + } + if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { + t.Fatalf("river state = %q, want cancelled", state) + } + time.Sleep(300 * time.Millisecond) + if n := runs.Load(); n != 0 { + t.Fatalf("cancelled queued job ran %d times", n) + } + }) + + t.Run("running", func(t *testing.T) { + id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Running", MaxAttempts: 1}) + if err != nil { + t.Fatal(err) + } + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("job did not start") + } + if err := m.CancelJob(ctx, id); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + waitRiverState(t, gdb, rec.RiverJobID, "cancelled") + rec = mustGet(t, m, id) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true (not ERROR)", rec.Status, rec.IsCanceled) + } + }) +} + +// TestQueueClear covers D-05: pending jobs of one queue are deleted and a +// running job is left alone. +func TestQueueClear(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, map[string]any{"queue.queues.default": 1}) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + release := make(chan struct{}) + started := make(chan struct{}, 1) + var pings atomic.Int32 + block := Job(func(ctx context.Context, a blockArgs) error { + started <- struct{}{} + select { + case <-release: + return nil + case <-ctx.Done(): + return ctx.Err() + } + }) + ping := Job(func(ctx context.Context, a pingArgs) error { + pings.Add(1) + return nil + }) + startTestWorker(t, app, WorkerOptions{}, block, ping) + ctx := t.Context() + if err := m.Enqueue(ctx, nil, blockArgs{}, EnqueueOpts{}); err != nil { + t.Fatal(err) + } + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("blocking job did not start") + } + for range 2 { + if err := m.Enqueue(ctx, nil, pingArgs{Tag: "pending"}, EnqueueOpts{}); err != nil { + t.Fatal(err) + } + } + + runClear := func() string { + t.Helper() + var out bytes.Buffer + root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), &out) + if err != nil { + t.Fatal(err) + } + root.SetArgs([]string{"queue:clear"}) + if err := root.ExecuteContext(ctx); err != nil { + t.Fatalf("queue:clear: %v\n%s", err, out.String()) + } + return out.String() + } + out := runClear() + if !strings.Contains(out, `Clearing queue "default"`) || !strings.Contains(out, "Cleared 2 jobs") { + t.Fatalf("output = %q", out) + } + close(release) + deadline := time.Now().Add(10 * time.Second) + for countRows(t, gdb, `SELECT count(*) FROM river_job WHERE kind = ? AND state = 'completed'`, "conga_test_block") != 1 { + if time.Now().After(deadline) { + t.Fatal("running job did not complete after queue:clear") + } + time.Sleep(25 * time.Millisecond) + } + if n := countRows(t, gdb, `SELECT count(*) FROM river_job WHERE kind = ?`, "conga_test_ping"); n != 0 || pings.Load() != 0 { + t.Fatalf("ping jobs left %d, ran %d", n, pings.Load()) + } + if out := runClear(); !strings.Contains(out, "Cleared 0 jobs") { + t.Fatalf("second clear output = %q", out) + } +} + +// syncBuffer is a bytes.Buffer safe for a command writing while the test +// reads. +type syncBuffer struct { + mu chan struct{} + buf bytes.Buffer +} + +func newSyncBuffer() *syncBuffer { return &syncBuffer{mu: make(chan struct{}, 1)} } + +func (b *syncBuffer) Write(p []byte) (int, error) { + b.mu <- struct{}{} + defer func() { <-b.mu }() + return b.buf.Write(p) +} + +func (b *syncBuffer) String() string { + b.mu <- struct{}{} + defer func() { <-b.mu }() + return b.buf.String() +} + +// TestQueueWorkCommand covers CLI-06: queue:work runs a filtered worker +// until its context ends and rejects unknown queues. +func TestQueueWorkCommand(t *testing.T) { + _, dsn := migratedDB(t) + + t.Run("unknown_queue", func(t *testing.T) { + var out bytes.Buffer + root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), &out) + if err != nil { + t.Fatal(err) + } + root.SetArgs([]string{"queue:work", "--queue", "nope"}) + err = root.ExecuteContext(t.Context()) + if err == nil || !strings.Contains(err.Error(), "known queues: default") { + t.Fatalf("err = %v, want unknown queue listing known queues", err) + } + }) + + t.Run("runs_until_cancelled", func(t *testing.T) { + out := newSyncBuffer() + root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), out) + if err != nil { + t.Fatal(err) + } + root.SetArgs([]string{"queue:work", "--queue", "default"}) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { done <- root.ExecuteContext(ctx) }() + deadline := time.Now().Add(10 * time.Second) + for !strings.Contains(out.String(), "worker started on queues: default") { + if time.Now().After(deadline) { + t.Fatalf("worker did not start: %q", out.String()) + } + time.Sleep(25 * time.Millisecond) + } + cancel() + select { + case err := <-done: + if err != nil { + t.Fatalf("queue:work returned %v", err) + } + case <-time.After(15 * time.Second): + t.Fatal("queue:work did not stop") + } + }) +} + +// TestStartServeWorker covers D-17: serve runs a worker unless +// queue.work_in_serve is false. +func TestStartServeWorker(t *testing.T) { + db, dsn := migratedDB(t) + off, _ := testApp(t, db, dsn, map[string]any{"queue.work_in_serve": false}) + w, err := StartServeWorker(t.Context(), off, nil) + if err != nil || w != nil { + t.Fatalf("work_in_serve false: worker %v, err %v", w, err) + } + on, _ := testApp(t, db, dsn, nil) + w, err = StartServeWorker(t.Context(), on, nil) + if err != nil || w == nil { + t.Fatalf("default: worker %v, err %v", w, err) + } + if got := w.Queues(); !reflect.DeepEqual(got, []string{"default"}) { + t.Fatalf("queues = %v", got) + } + if err := w.Stop(context.Background()); err != nil { + t.Fatal(err) + } +} diff --git a/modules/conga/worker.go b/modules/conga/worker.go index 4f57aa0..1aa2464 100644 --- a/modules/conga/worker.go +++ b/modules/conga/worker.go @@ -28,6 +28,8 @@ type WorkerOptions struct { pollInterval time.Duration // pollOnly builds a driver without the LISTEN pool (tests only). pollOnly bool + // retryPolicy overrides River's retry schedule (tests only). + retryPolicy river.ClientRetryPolicy } // Worker is a running River work client. Stop it with Stop. @@ -92,6 +94,13 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, return nil, fmt.Errorf("conga: register %s: %w", j.kind(), err) } } + if len(m.jobs) == 0 { + // River refuses to start a client without workers; an app with no + // jobs yet still gets a worker that starts and idles. + if err := river.AddWorkerSafely[idleArgs](workers, idleWorker{}); err != nil { + return nil, fmt.Errorf("conga: register idle worker: %w", err) + } + } cfg := baseConfig(s, log) cfg.Workers = workers cfg.Queues = map[string]river.QueueConfig{} @@ -101,6 +110,9 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, if o.pollInterval > 0 { cfg.FetchPollInterval = o.pollInterval } + if o.retryPolicy != nil { + cfg.RetryPolicy = o.retryPolicy + } var listener *pgxpool.Pool var driver *riverdatabasesql.Driver @@ -126,6 +138,16 @@ func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, return &Worker{m: m, client: client, listener: listener, queues: sortedQueueNames(queues)}, nil } +// StartServeWorker starts the in-process worker of the serve command on +// every known queue. It returns nil, nil when queue.work_in_serve is false, +// for deployments that run a separate queue:work process. +func StartServeWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin) (*Worker, error) { + if !settingsFromApp(app).workInServe { + return nil, nil + } + return StartWorker(ctx, app, plugins, WorkerOptions{}) +} + // Stop stops fetching jobs and waits for running jobs to finish. When ctx // ends first, running jobs are cancelled. The listener pool is closed. func (w *Worker) Stop(ctx context.Context) error { @@ -242,3 +264,15 @@ func summerJobID(row *rivertype.JobRow) (uint, bool) { } return meta.SummerJobID, true } + +// idleArgs is the kind of the placeholder worker registered when an app has +// no jobs. Nothing inserts it. +type idleArgs struct{} + +func (idleArgs) Kind() string { return "summercms_conga_idle" } + +type idleWorker struct { + river.WorkerDefaults[idleArgs] +} + +func (idleWorker) Work(context.Context, *river.Job[idleArgs]) error { return nil } diff --git a/modules/surf/README.md b/modules/surf/README.md index ce51926..d6eaf2d 100644 --- a/modules/surf/README.md +++ b/modules/surf/README.md @@ -151,18 +151,18 @@ http: supports_credentials: true ``` -The `serve` command also opens the database through [lagoon](../lagoon/README.md) and the uploads bucket through `attach.OpenBucket`, so their settings must be present as well. +The `serve` command also opens the database through [lagoon](../lagoon/README.md) and the uploads bucket through `attach.OpenBucket`, so their settings must be present as well. It starts the background job worker of [conga](../conga/README.md) in the same process unless `queue.work_in_serve` is `false` (default `true`); set it to `false` when a separate `queue:work` process runs the jobs. The other `queue.*` keys are documented in the conga README. ## CLI commands | Command | Flags | Description | |---------|-------|-------------| -| `serve` | `--addr` (default `:8080`) | Opens the database and uploads bucket, assembles the router and serves HTTP until SIGINT or SIGTERM, then shuts down gracefully within 10 seconds. | +| `serve` | `--addr` (default `:8080`) | Opens the database and uploads bucket, assembles the router, starts the in-process job worker (see `queue.work_in_serve`) and serves HTTP until SIGINT or SIGTERM, then shuts the server and the worker down gracefully within 10 seconds. | | `route:list` | none | Builds the router the same way `serve` does, without opening the database or listening, and prints a table of method, pattern, plugin, middleware and raw flag for every route. | ## Dependencies -- SummerCMS modules: [backpack](../backpack/README.md), [bonfire](../bonfire/README.md), [bouncer](../bouncer/README.md), [cabana](../cabana/README.md), [compass](../compass/README.md), [lagoon](../lagoon/README.md) (including `lagoon/attach`), [pact](../pact/README.md), [party](../party/README.md), [towel](../towel/README.md), [wire](../wire/README.md). +- SummerCMS modules: [backpack](../backpack/README.md), [bonfire](../bonfire/README.md), [bouncer](../bouncer/README.md), [cabana](../cabana/README.md), [compass](../compass/README.md), [conga](../conga/README.md) (the in-process job worker of `serve`), [lagoon](../lagoon/README.md) (including `lagoon/attach`), [pact](../pact/README.md), [party](../party/README.md), [towel](../towel/README.md), [wire](../wire/README.md). - Third-party: `gocloud.dev/blob` (uploads bucket opened by `serve`). - Standard library: `bytes`, `context`, `fmt`, `math`, `net`, `net/http`, `net/netip`, `os`, `os/signal`, `regexp`, `strconv`, `strings`, `sync`, `syscall`, `time`. diff --git a/modules/surf/serve.go b/modules/surf/serve.go index 02f2c03..ff67482 100644 --- a/modules/surf/serve.go +++ b/modules/surf/serve.go @@ -13,6 +13,7 @@ import ( "git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/bonfire" + "git.golem15.com/golem15/summercms/modules/conga" "git.golem15.com/golem15/summercms/modules/lagoon" "git.golem15.com/golem15/summercms/modules/lagoon/attach" "git.golem15.com/golem15/summercms/modules/party" @@ -34,6 +35,8 @@ func ServeCommand(app *backpack.App, plugins []party.Plugin) bonfire.Command { if strings.TrimSpace(addr) == "" { addr = ":8080" } + ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) + defer stop() sqlDB, gdb, err := lagoon.OpenFromApp(ctx, app) if err != nil { return err @@ -51,36 +54,57 @@ func ServeCommand(app *backpack.App, plugins []party.Plugin) bonfire.Command { if err != nil { return err } - ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) - defer stop() + worker, err := conga.StartServeWorker(ctx, app, plugins) + if err != nil { + return err + } + if worker != nil { + out.Info(fmt.Sprintf("job worker started on queues: %s", strings.Join(worker.Queues(), ", "))) + } srv := &http.Server{Addr: addr, Handler: h, BaseContext: func(net.Listener) context.Context { return ctx }} errCh := make(chan error, 1) go func() { out.Info(fmt.Sprintf("listening on %s", addr)) errCh <- srv.ListenAndServe() }() + shutdownCtx := func() (context.Context, context.CancelFunc) { + return context.WithTimeout(context.Background(), 10*time.Second) + } select { case <-ctx.Done(): - shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + sctx, cancel := shutdownCtx() defer cancel() - if err := srv.Shutdown(shutdownCtx); err != nil { - return err + shutdownErr := srv.Shutdown(sctx) + workerErr := stopWorker(sctx, worker) + if shutdownErr != nil { + return shutdownErr } err := <-errCh - if err == nil || err == http.ErrServerClosed { - return nil + if err != nil && err != http.ErrServerClosed { + return err } - return err + return workerErr case err := <-errCh: - if err == nil || err == http.ErrServerClosed { - return nil + sctx, cancel := shutdownCtx() + defer cancel() + workerErr := stopWorker(sctx, worker) + if err != nil && err != http.ErrServerClosed { + return err } - return err + return workerErr } }, } } +// stopWorker stops the in-process job worker, if serve started one. +func stopWorker(ctx context.Context, w *conga.Worker) error { + if w == nil { + return nil + } + return w.Stop(ctx) +} + // publishUploads opens storage.uploads.bucket_url and publishes *blob.Bucket. // An empty URL fails boot the same way an empty JWT secret does. func publishUploads(ctx context.Context, app *backpack.App) (*blob.Bucket, error) {