Files
summercms/modules/conga/example_test.go
Jakub Zych 9d37d56486 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
2026-09-30 22:35:22 +02:00

177 lines
5.1 KiB
Go

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)
}
}