- 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
177 lines
5.1 KiB
Go
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)
|
|
}
|
|
}
|