From 7457c44be5ec4f9247154d2fdec2fe8885db3f45 Mon Sep 17 00:00:00 2001 From: Robin Jarry Date: Sat, 14 Dec 2024 23:53:26 +0100 Subject: [PATCH] treewide: update to new multi worker dowork api The work.Queue() implementation has changed and now allows scheduling tasks from multiple goroutines in parallel. Also, it requires a new argument for limiting the queue buffer size in order to apply back pressure. Define new configuration variables to allow tuning the queue sizes and number of workers per service: [mail] # Maximum size of the outgoing email queue (default 512). egress-queue-size = 512 [$service_name] # Number of parallel workers per queue (default 1). # There are multiple queues (for egress email, webhooks, etc.). # This setting is applied on a per-queue basis. queue-workers = 1 [webhooks] # Maximum size of the webhooks queue (default 512). queue-size = 512 Use these new settings to configure the queues and workers accordingly. Fix unit tests as well. Link: https://git.sr.ht/~sircmpwn/dowork/commit/95719cfc0118 Signed-off-by: Robin Jarry --- email/worker.go | 11 ++++++++++- go.mod | 2 +- go.sum | 4 ++-- server/server.go | 24 +++++++++++++++-------- webhooks/legacy.go | 13 +++++++++++-- webhooks/legacy_test.go | 42 ++++++++++++++++++----------------------- webhooks/queue.go | 13 +++++++++++-- 7 files changed, 69 insertions(+), 40 deletions(-) diff --git a/email/worker.go b/email/worker.go index 4f3501b5b0bf94d7f28486944601ff3b1cd7388c..7c75478bc4cc7e747ae7b8467a00746e2081c954 100644 --- a/email/worker.go +++ b/email/worker.go @@ -9,6 +9,7 @@ import ( "log" "net/http" "os" + "strconv" "strings" "time" @@ -228,8 +229,16 @@ func NewQueue(conf ini.File) *Queue { } } + queueSize := 512 + if s, ok := conf.Get("mail", "egress-queue-size"); ok { + var err error + if queueSize, err = strconv.Atoi(s); err != nil { + panic(fmt.Errorf("[mail]egress-queue-size: %w", err)) + } + } + return &Queue{ - Queue: work.NewQueue("email"), + Queue: work.NewQueue("email", queueSize), smtpFrom: addr, ownerAddress: ownerAddr, entity: entity, diff --git a/go.mod b/go.mod index 1db00edfe0c05b002e05ce44f23e71f1bf8a453d..d902a16ee449bec37f0cb4c953d140286502e694 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,7 @@ module git.sr.ht/~sircmpwn/core-go go 1.22 require ( - git.sr.ht/~sircmpwn/dowork v0.0.0-20221010085743-46c4299d76a1 + git.sr.ht/~sircmpwn/dowork v0.0.0-20241209140539-95719cfc0118 git.sr.ht/~sircmpwn/getopt v1.0.0 git.sr.ht/~sircmpwn/go-bare v0.0.0-20210406120253-ab86bc2846d9 github.com/99designs/gqlgen v0.17.36 diff --git a/go.sum b/go.sum index 1a1074ea51358e91cc96b416e6ea3c686fb144f1..4d137d63f2472cbf43e00b353ec9c671711c308c 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,5 @@ -git.sr.ht/~sircmpwn/dowork v0.0.0-20221010085743-46c4299d76a1 h1:EvPKkneKkF/f7zEgKPqIZVyj3jWO8zSmsBOvMhAGqMA= -git.sr.ht/~sircmpwn/dowork v0.0.0-20221010085743-46c4299d76a1/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= +git.sr.ht/~sircmpwn/dowork v0.0.0-20241209140539-95719cfc0118 h1:aKZ3Es4Tj6IKLGilwcPlLzeXwrz+/ZalVimyiX6bmdM= +git.sr.ht/~sircmpwn/dowork v0.0.0-20241209140539-95719cfc0118/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= git.sr.ht/~sircmpwn/getopt v0.0.0-20191230200459-23622cc906b3/go.mod h1:wMEGFFFNuPos7vHmWXfszqImLppbc0wEhh6JBfJIUgw= git.sr.ht/~sircmpwn/getopt v1.0.0 h1:/pRHjO6/OCbBF4puqD98n6xtPEgE//oq5U8NXjP7ROc= git.sr.ht/~sircmpwn/getopt v1.0.0/go.mod h1:wMEGFFFNuPos7vHmWXfszqImLppbc0wEhh6JBfJIUgw= diff --git a/server/server.go b/server/server.go index 92547f9b8109a756c0765f9bdbd42f2f35be5a5d..62e35a38181e1d5a8c055c661b8496ca0bc81ce5 100644 --- a/server/server.go +++ b/server/server.go @@ -260,16 +260,24 @@ func (server *Server) WithMiddleware( // Add dowork task queues for this server to manage func (server *Server) WithQueues(queues ...*work.Queue) *Server { - ctx := context.Background() - ctx = config.Context(ctx, server.conf, server.service) - ctx = database.Context(ctx, server.db) - ctx = redis.Context(ctx, server.redis) - ctx = email.Context(ctx, server.email) - ctx = context.WithValue(ctx, serverCtxKey, server) - + queueWorkers := 1 + if n, ok := server.conf.Get(server.service, "queue-workers"); ok { + var err error + if queueWorkers, err = strconv.Atoi(n); err != nil { + panic(fmt.Errorf("[%s]queue-workers: %w", server.service, err)) + } + } server.queues = append(server.queues, queues...) for _, queue := range queues { - queue.Start(ctx) + // Use a different context per worker to allow "goroutine-local" + // variables. + ctx := context.Background() + ctx = config.Context(ctx, server.conf, server.service) + ctx = database.Context(ctx, server.db) + ctx = redis.Context(ctx, server.redis) + ctx = email.Context(ctx, server.email) + ctx = context.WithValue(ctx, serverCtxKey, server) + queue.Start(ctx, queueWorkers) } return server } diff --git a/webhooks/legacy.go b/webhooks/legacy.go index a037c7e8a1864821a58697ed2a5d43a396ad767f..d3838d64d9673fe7be5e39095e136b1e9c729fb6 100644 --- a/webhooks/legacy.go +++ b/webhooks/legacy.go @@ -9,12 +9,14 @@ import ( "io/ioutil" "log" "net/http" + "strconv" "strings" "time" "git.sr.ht/~sircmpwn/dowork" sq "github.com/Masterminds/squirrel" "github.com/google/uuid" + "github.com/vaughan0/go-ini" "git.sr.ht/~sircmpwn/core-go/crypto" "git.sr.ht/~sircmpwn/core-go/database" @@ -33,9 +35,16 @@ type LegacySubscription struct { // Creates a new worker for delivering legacy webhooks. The caller must start // the worker themselves. -func NewLegacyQueue() *LegacyQueue { +func NewLegacyQueue(conf ini.File) *LegacyQueue { + queueSize := 512 + if s, ok := conf.Get("webhooks", "queue-size"); ok { + var err error + if queueSize, err = strconv.Atoi(s); err != nil { + panic(fmt.Errorf("[webhooks]queue-size: %w", err)) + } + } return &LegacyQueue{ - work.NewQueue("webhooks_legacy"), + work.NewQueue("webhooks_legacy", queueSize), } } diff --git a/webhooks/legacy_test.go b/webhooks/legacy_test.go index a0b6d65af6e83d83ef2465889d43be5b76129576..b3bfb685dd17f8a226c118b9b41a231860530ff5 100644 --- a/webhooks/legacy_test.go +++ b/webhooks/legacy_test.go @@ -57,12 +57,11 @@ func (ac *argContains) Match(v driver.Value) bool { } func TestDelivery(t *testing.T) { - var called bool + called := make(chan struct{}) srv := httptest.NewServer(http.HandlerFunc( func(w http.ResponseWriter, r *http.Request) { defer r.Body.Close() - called = true assert.Equal(t, r.Method, http.MethodPost) assert.Equal(t, r.URL.Path, "/webhook") @@ -79,6 +78,7 @@ func TestDelivery(t *testing.T) { assert.True(t, crypto.VerifyWebhook(b, nonce, signature)) w.Write([]byte("Thanks!")) + close(called) })) defer srv.Close() @@ -87,6 +87,8 @@ func TestDelivery(t *testing.T) { panic(err) } ctx := database.Context(context.Background(), db) + ctx, cancel := context.WithTimeout(ctx, 1*time.Second) + defer cancel() // Lookup phase mock.ExpectBegin() @@ -97,30 +99,12 @@ func TestDelivery(t *testing.T) { srv.URL+"/webhook", "profile:update")). WithArgs(42, sqlmock.AnyArg()) // Any => events LIKE %profile:update% mock.ExpectCommit() - - queue := NewLegacyQueue() - q := sq. - Select(). - From("user_webhook_subscription sub"). - Where(`sub.user_id = ?`, 42) - queue.Schedule(ctx, q, "user", "profile:update", []byte(`{"hello": "world"}`)) - // Schedule phase mock.ExpectBegin() mock.ExpectQuery(`INSERT INTO user_webhook_delivery`). WillReturnRows(sqlmock.NewRows([]string{"id"}).AddRow(4096)) mock.ExpectCommit() - - queue.Queue.Dispatch(ctx) - - assert.Nil(t, mock.ExpectationsWereMet()) - // Delivery phase - db, mock, err = sqlmock.New() - if err != nil { - panic(err) - } - mock.ExpectBegin() mock.ExpectExec(`UPDATE user_webhook_delivery`). WithArgs("Thanks!", 200, @@ -135,9 +119,19 @@ func TestDelivery(t *testing.T) { WillReturnResult(sqlmock.NewResult(1, 1)) mock.ExpectCommit() - ctx = database.Context(context.Background(), db) - queue.Queue.Dispatch(ctx) + queue := NewLegacyQueue(make(ini.File)) + queue.Queue.Start(ctx, 1) + q := sq. + Select(). + From("user_webhook_subscription sub"). + Where(`sub.user_id = ?`, 42) + queue.Schedule(ctx, q, "user", "profile:update", []byte(`{"hello": "world"}`)) - assert.Nil(t, mock.ExpectationsWereMet()) - assert.True(t, called) + select { + case <-called: + queue.Queue.Shutdown() + assert.Nil(t, mock.ExpectationsWereMet()) + case <-ctx.Done(): + t.Fatal("webhook url endpoint not called") + } } diff --git a/webhooks/queue.go b/webhooks/queue.go index 8286034d4d2340d8704fb5fd694dcdb977dda652..172043fb23f152da3c1b4914deb1be1cb155e526 100644 --- a/webhooks/queue.go +++ b/webhooks/queue.go @@ -9,6 +9,7 @@ import ( "io/ioutil" "log" "net/http" + "strconv" "strings" "time" @@ -16,6 +17,7 @@ import ( "github.com/99designs/gqlgen/graphql" sq "github.com/Masterminds/squirrel" "github.com/google/uuid" + "github.com/vaughan0/go-ini" "git.sr.ht/~sircmpwn/core-go/auth" "git.sr.ht/~sircmpwn/core-go/crypto" @@ -42,8 +44,15 @@ type WebhookSubscription struct { // Creates a new worker for delivering webhooks. The caller must start the // worker themselves. -func NewQueue(schema graphql.ExecutableSchema) *WebhookQueue { - return &WebhookQueue{work.NewQueue("webhooks"), schema} +func NewQueue(schema graphql.ExecutableSchema, conf ini.File) *WebhookQueue { + queueSize := 512 + if s, ok := conf.Get("webhooks", "queue-size"); ok { + var err error + if queueSize, err = strconv.Atoi(s); err != nil { + panic(fmt.Errorf("[webhooks]queue-size: %w", err)) + } + } + return &WebhookQueue{work.NewQueue("webhooks", queueSize), schema} } // Schedules delivery of a webhook to a set of subscribers.