From af2afebd4c0a5b082dedc2a21aa00a57a65c65ca Mon Sep 17 00:00:00 2001 From: Drew DeVault Date: Sat, 10 Oct 2020 11:24:59 -0400 Subject: [PATCH] Add legacy webhooks worker implementation --- auth/middleware.go | 8 +- auth/middleware_test.go | 10 +- crypto/crypto.go | 10 ++ crypto/crypto_test.go | 6 + database/sq.go | 4 +- go.mod | 5 +- go.sum | 7 ++ server/server.go | 22 ++-- webhooks/legacy.go | 250 ++++++++++++++++++++++++++++++++++++++++ webhooks/legacy_test.go | 114 ++++++++++++++++++ 10 files changed, 412 insertions(+), 24 deletions(-) create mode 100644 webhooks/legacy.go create mode 100644 webhooks/legacy_test.go diff --git a/auth/middleware.go b/auth/middleware.go index 3e5b3252e0cb174589a77ae3909284317384f8d1..3383fcb4ff485ea135dde0a120d84f637bc039a9 100644 --- a/auth/middleware.go +++ b/auth/middleware.go @@ -91,7 +91,7 @@ func authForUsername(ctx context.Context, username string) (*AuthContext, error) var auth AuthContext if err := database.WithTx(ctx, &sql.TxOptions{ Isolation: 0, - ReadOnly: true, + ReadOnly: true, }, func(tx *sql.Tx) error { var ( err error @@ -149,7 +149,7 @@ func authForOAuthClient(ctx context.Context, clientUUID string) (*AuthContext, e var auth AuthContext if err := database.WithTx(ctx, &sql.TxOptions{ Isolation: 0, - ReadOnly: true, + ReadOnly: true, }, func(tx *sql.Tx) error { var ( err error @@ -399,7 +399,7 @@ func FetchMetaProfile(ctx context.Context, username string, user *AuthContext) e func LookupUser(ctx context.Context, username string, user *AuthContext) error { return database.WithTx(ctx, &sql.TxOptions{ Isolation: 0, - ReadOnly: true, + ReadOnly: true, }, func(tx *sql.Tx) error { var ( err error @@ -578,7 +578,7 @@ func LegacyOAuth(bearer string, hash [64]byte, w http.ResponseWriter, ) if err := database.WithTx(r.Context(), &sql.TxOptions{ Isolation: 0, - ReadOnly: true, + ReadOnly: true, }, func(tx *sql.Tx) error { var ( err error diff --git a/auth/middleware_test.go b/auth/middleware_test.go index 2ba8e89c0faa3075c10fc4f7d3528d47e38bcf54..59fbe0c1d993398a184385f4b5b9cd8fe5e6ac4d 100644 --- a/auth/middleware_test.go +++ b/auth/middleware_test.go @@ -93,14 +93,14 @@ func TestInternal(t *testing.T) { assert.Nil(t, err) req.Header.Add("Content-Type", "application/json") internalAuth := InternalAuth{ - Name: "jdoe", - ClientID: "", - NodeID: "test.node", + Name: "jdoe", + ClientID: "", + NodeID: "test.node", OAuthClientUUID: "", } payload, err := json.Marshal(&internalAuth) assert.Nil(t, err) - req.Header.Add("Authorization", "Internal " + + req.Header.Add("Authorization", "Internal "+ string(crypto.Encrypt(payload))) req.RemoteAddr = "127.0.0.1" @@ -122,7 +122,7 @@ func TestInternal(t *testing.T) { strings.NewReader(`{"query": "query { me { id } }"}`)) assert.Nil(t, err) req.Header.Add("Content-Type", "application/json") - req.Header.Add("Authorization", "Internal " + + req.Header.Add("Authorization", "Internal "+ string(crypto.Encrypt(payload))) req.RemoteAddr = "1.2.3.4" diff --git a/crypto/crypto.go b/crypto/crypto.go index 7e28e7dcbb30f09d2c69748031d78564d3793a1f..025d7453ae5cbbf1bb183ddac277c9d00e5ff874 100644 --- a/crypto/crypto.go +++ b/crypto/crypto.go @@ -85,6 +85,8 @@ func HMACVerify(payload []byte, signature []byte) bool { return hmac.Equal(expected, signature) } +// Signs the payload for a webhook, returning respectively the values for the +// X-Payload-Nonce and X-Payload-Signature headers. func SignWebhook(payload []byte) (string, string) { var nonceSeed [8]byte _, err := rand.Read(nonceSeed[:]) @@ -97,3 +99,11 @@ func SignWebhook(payload []byte) (string, string) { Sign(append(payload, []byte(nonce)...))) return nonce, signature } + +func VerifyWebhook(payload []byte, nonce, signature string) bool { + s, err := base64.StdEncoding.DecodeString(signature) + if err != nil { + return false + } + return Verify(append(payload, []byte(nonce)...), s) +} diff --git a/crypto/crypto_test.go b/crypto/crypto_test.go index 81d872bf47fbf8216c5f280692d6b472af176dc0..eb7e7cad4519715e79e7310690f7259881668b57 100644 --- a/crypto/crypto_test.go +++ b/crypto/crypto_test.go @@ -33,6 +33,12 @@ func TestSignWebhook(t *testing.T) { assert.True(t, valid) } +func TestVerifyWebhook(t *testing.T) { + payload := []byte("Hello world!") + nonce, signature := SignWebhook(payload) + assert.True(t, VerifyWebhook(payload, nonce, signature)) +} + func TestSign(t *testing.T) { payload := []byte("Hello world!") signature := Sign(payload) diff --git a/database/sq.go b/database/sq.go index 58b558e0f3b4b9d7ca9b7f07489c42cd2ccb0780..3fbe5ab9d5da06719e9135f18de8a3cec5adc60a 100644 --- a/database/sq.go +++ b/database/sq.go @@ -69,9 +69,9 @@ func (mf *ModelFields) Anonymous() []*FieldMap { } type Model interface { - Alias() string + Alias() string Fields() *ModelFields - Table() string + Table() string } func Select(ctx context.Context, cols ...interface{}) sq.SelectBuilder { diff --git a/go.mod b/go.mod index 84e3759137a181b9ac2d410f6132a56a28487ab5..5ea29190fd0db1aae31b1808d7cf791cce4807ff 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,7 @@ module git.sr.ht/~sircmpwn/core-go go 1.13 require ( - git.sr.ht/~sircmpwn/dowork v0.0.0-20201006201820-f2599e406ecb + git.sr.ht/~sircmpwn/dowork v0.0.0-20201009195117-58f54dd4c3e8 git.sr.ht/~sircmpwn/getopt v0.0.0-20191230200459-23622cc906b3 git.sr.ht/~sircmpwn/go-bare v0.0.0-20200812160916-d2c72e1a5018 github.com/99designs/gqlgen v0.13.0 @@ -12,6 +12,7 @@ require ( github.com/fernet/fernet-go v0.0.0-20191111064656-eff2850e6001 github.com/go-chi/chi v4.1.2+incompatible github.com/go-redis/redis/v8 v8.2.3 + github.com/google/uuid v1.0.0 github.com/kavu/go_reuseport v1.5.0 github.com/lib/pq v1.8.0 github.com/martinlindhe/base36 v1.1.0 @@ -23,7 +24,7 @@ require ( github.com/vektah/gqlparser v1.3.1 github.com/vektah/gqlparser/v2 v2.1.0 golang.org/x/crypto v0.0.0-20200728195943-123391ffb6de - golang.org/x/sys v0.0.0-20201007082116-8445cc04cbdf // indirect + golang.org/x/sys v0.0.0-20201009025420-dfb3f7c4e634 // indirect google.golang.org/protobuf v1.25.0 // indirect gopkg.in/mail.v2 v2.3.1 ) diff --git a/go.sum b/go.sum index b2f4f99b7ec70e731c9d8404471bd91057d7c0d9..dd30b17b51e9d07b3a39888b7f290c9be17ca63d 100644 --- a/go.sum +++ b/go.sum @@ -4,6 +4,10 @@ git.sr.ht/~sircmpwn/dowork v0.0.0-20201002192337-cc78e95c493c h1:DHYVIt2TT6Nx+CK git.sr.ht/~sircmpwn/dowork v0.0.0-20201002192337-cc78e95c493c/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= git.sr.ht/~sircmpwn/dowork v0.0.0-20201006201820-f2599e406ecb h1:2Yodrugga89JpSewI+TOj4zUOeujnDdpIAvC91qIpys= git.sr.ht/~sircmpwn/dowork v0.0.0-20201006201820-f2599e406ecb/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= +git.sr.ht/~sircmpwn/dowork v0.0.0-20201009194917-181b76f9491b h1:7tYMLNLFAjN3R+Q2mRGjQ4TDbYHUdoMmqlTB9++gTB0= +git.sr.ht/~sircmpwn/dowork v0.0.0-20201009194917-181b76f9491b/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= +git.sr.ht/~sircmpwn/dowork v0.0.0-20201009195117-58f54dd4c3e8 h1:MSiW/2sDb2KjidhG/orqaQZfu0bpQcKr9Sbu8aRqS9o= +git.sr.ht/~sircmpwn/dowork v0.0.0-20201009195117-58f54dd4c3e8/go.mod h1:8neHEO3503w/rNtttnR0JFpQgM/GFhaafVwvkPsFIDw= git.sr.ht/~sircmpwn/getopt v0.0.0-20191230200459-23622cc906b3 h1:4wDp4BKF7NQqoh73VXpZsB/t1OEhDpz/zEpmdQfbjDk= git.sr.ht/~sircmpwn/getopt v0.0.0-20191230200459-23622cc906b3/go.mod h1:wMEGFFFNuPos7vHmWXfszqImLppbc0wEhh6JBfJIUgw= git.sr.ht/~sircmpwn/go-bare v0.0.0-20200812160916-d2c72e1a5018 h1:89QMorzx6ML69PKPoayL3HuSfb7WqAlxD1dZ7DyzD0k= @@ -125,6 +129,7 @@ github.com/google/go-cmp v0.5.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/ github.com/google/go-cmp v0.5.1/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/renameio v0.1.0/go.mod h1:KWCgfxg9yswjAJkECMjeO8J8rahYeXnNhOm40UhjYkI= +github.com/google/uuid v1.0.0 h1:b4Gk+7WdP/d3HZH8EJsZpvV7EtDOgaZLtnaNGIu1adA= github.com/google/uuid v1.0.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1/go.mod h1:wJfORRmW1u3UXTncJ5qlYoELFm8eSnnEO6hX4iZ3EWY= github.com/gorilla/context v0.0.0-20160226214623-1ea25387ff6f/go.mod h1:kBGZzfjB9CEq2AlWe17Uuf7NDRt0dE0s8S51q0aT7Yg= @@ -421,6 +426,8 @@ golang.org/x/sys v0.0.0-20200615200032-f1bc736245b1/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20200625212154-ddb9806d33ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201007082116-8445cc04cbdf h1:AvBTl0xbF/KtHyvm61X4gSPF7/dKJ/xQqJwKr1Qu9no= golang.org/x/sys v0.0.0-20201007082116-8445cc04cbdf/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20201009025420-dfb3f7c4e634 h1:bNEHhJCnrwMKNMmOx3yAynp5vs5/gRy+XWFtZFu7NBM= +golang.org/x/sys v0.0.0-20201009025420-dfb3f7c4e634/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= golang.org/x/time v0.0.0-20180412165947-fbb02b2291d2/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= diff --git a/server/server.go b/server/server.go index e38dafebe683375e89de9c11858c2ca42ee4afe5..79c1d5a74c945ccde07539ca6d364ce6a94c7013 100644 --- a/server/server.go +++ b/server/server.go @@ -19,13 +19,13 @@ import ( "github.com/99designs/gqlgen/handler" "github.com/go-chi/chi" "github.com/go-chi/chi/middleware" + goRedis "github.com/go-redis/redis/v8" "github.com/kavu/go_reuseport" + _ "github.com/lib/pq" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promauto" "github.com/prometheus/client_golang/prometheus/promhttp" "github.com/vaughan0/go-ini" - _ "github.com/lib/pq" - goRedis "github.com/go-redis/redis/v8" "git.sr.ht/~sircmpwn/core-go/auth" "git.sr.ht/~sircmpwn/core-go/config" @@ -46,13 +46,13 @@ var ( ) type Server struct { - conf ini.File - db *sql.DB - redis *goRedis.Client - router chi.Router - schema graphql.ExecutableSchema - service string - queues []*work.Queue + conf ini.File + db *sql.DB + redis *goRedis.Client + router chi.Router + schema graphql.ExecutableSchema + service string + queues []*work.Queue } // Creates a new common server context for a SourceHut GraphQL daemon. @@ -83,7 +83,7 @@ func (server *Server) WithSchema( err error ) if limit, ok := server.conf.Get( - server.service + "::api", "max-complexity"); ok { + server.service+"::api", "max-complexity"); ok { complexity, err = strconv.Atoi(limit) if err != nil { panic(err) @@ -230,7 +230,7 @@ func (server *Server) Run() { log.Println("Terminating server...") ctx, cancel := context.WithDeadline(context.Background(), - time.Now().Add(30 * time.Second)) + time.Now().Add(30*time.Second)) qserver.Shutdown(ctx) cancel() diff --git a/webhooks/legacy.go b/webhooks/legacy.go new file mode 100644 index 0000000000000000000000000000000000000000..4e431e1f75d8b7705d080316f5a7acf6ade58e18 --- /dev/null +++ b/webhooks/legacy.go @@ -0,0 +1,250 @@ +package webhooks + +import ( + "bytes" + "context" + "database/sql" + "fmt" + "io" + "io/ioutil" + "log" + "net/http" + "strings" + "time" + + "git.sr.ht/~sircmpwn/dowork" + sq "github.com/Masterminds/squirrel" + "github.com/google/uuid" + + "git.sr.ht/~sircmpwn/core-go/crypto" + "git.sr.ht/~sircmpwn/core-go/database" +) + +type LegacyQueue struct { + Queue *work.Queue +} + +type LegacySubscription struct { + ID int + Created time.Time + URL string + Events []string +} + +// Creates a new worker for delivering legacy webhooks. The caller must start +// the worker themselves. +func NewLegacyQueue() *LegacyQueue { + return &LegacyQueue{ + work.NewQueue("webhooks_legacy"), + } +} + +// Schedules delivery of a legacy webhook to a set of subscribers. +// +// The select builder should not return any columns, i.e. the caller should use +// squirrel.Select() with no parameters. The caller should prepare FROM and any +// WHERE clauses which are necessary to refine the subscriber list (e.g. by +// affected resource ID). +// +// Name shall be the prefix of the webhook tables, e.g. "user" for +// "user_webhook_{delivery,subscription}". +func (lq *LegacyQueue) Schedule(q sq.SelectBuilder, + name, event string, payload []byte) { + // The following tasks are done during this process: + // + // 1. Fetch subscription details from the database + // 2. Prepare deliveries and create delivery records + // 3. Deliver the webhooks + // + // The first two steps are done in this task, then N tasks are created for + // step 3 where N = number of subscriptions. + task := work.NewTask(func(ctx context.Context) error { + subs, err := fetchSubscriptions(ctx, q, event) + if err != nil { + return err + } + + if len(subs) == 0 { + return nil + } + + tasks := make([]*work.Task, len(subs)) + if err := database.WithTx(ctx, nil, func(tx *sql.Tx) error { + var err error + for i, sub := range subs { + tasks[i], err = lq.queueStage2(ctx, tx, + name, event, sub, payload) + if err != nil { + return err + } + } + return nil + }); err != nil { + return err + } + + for _, task := range tasks { + lq.Queue.Enqueue(task) + } + log.Printf("Enqueued %s %s webhook delivery for %d subscriptions", + name, event, len(subs)) + return nil + }) + lq.Queue.Enqueue(task) +} + +func fetchSubscriptions(ctx context.Context, q sq.SelectBuilder, + event string) ([]*LegacySubscription, error) { + + var subs []*LegacySubscription + if err := database.WithTx(ctx, &sql.TxOptions{ + Isolation: 0, + ReadOnly: true, + }, func(tx *sql.Tx) error { + var ( + err error + rows *sql.Rows + ) + if rows, err = q. + Columns("id", "created", "url", "events"). + Where(sq.Like{"events": "%" + event + "%"}). + RunWith(tx). + QueryContext(ctx); err != nil { + panic(err) + } + defer rows.Close() + + var events string + for rows.Next() { + var sub LegacySubscription + if err := rows.Scan(&sub.ID, &sub.Created, + &sub.URL, &events); err != nil { + panic(err) + } + + // The LIKE clause gets us an approximate list of implicated + // subscriptions, so we quickly decode the event list and + // double check here to get the final list. + sub.Events = strings.Split(events, ",") + + var valid bool + for _, e := range sub.Events { + if e == event { + valid = true + break + } + } + + if valid { + subs = append(subs, &sub) + } + } + return nil + }); err != nil { + return nil, err + } + + return subs, nil +} + +// Inserts the delivery record and schedules the actual delivery task +func (lq *LegacyQueue) queueStage2(ctx context.Context, tx *sql.Tx, + name, event string, sub *LegacySubscription, + payload []byte) (*work.Task, error) { + + deliveryUUID := uuid.New().String() + headers := make(http.Header) + headers.Set("Content-Type", "application/json") + headers.Set("X-Webhook-Event", event) + headers.Set("X-Webhook-Delivery", deliveryUUID) + var sb strings.Builder + headers.Write(&sb) + + var deliveryID int + sq.Insert(name+"_webhook_delivery"). + Columns("uuid", "created", "event", "url", + "payload", "payload_headers", "response_status", + "subscription_id"). + Values(deliveryUUID, "NOW() at time zone 'utc'", event, sub.URL, + string(payload), sb.String(), -2, sub.ID). + Suffix(`RETURNING (id)`). + RunWith(tx). + ScanContext(ctx, &deliveryID) + + return work.NewTask(func(ctx context.Context) error { + return deliverPayload(ctx, name, sub.URL, headers, payload, deliveryID) + }).Retries(5).After(func(ctx context.Context, task *work.Task) { + if task.Result() == nil { + log.Printf("LEGACY WEBHOOK: %s: delivery complete after %d attempts", + deliveryUUID, task.Attempts()) + } else { + log.Printf("LEGACY WEBHOOK: %s: delivery failed after %d attempts: %v", + deliveryUUID, task.Attempts(), task.Result()) + } + }), nil +} + +// Performs a webhook delivery and updates the delivery record in the database +func deliverPayload(ctx context.Context, name, url string, + headers http.Header, payload []byte, deliveryID int) error { + + client := &http.Client{ + Timeout: 30 * time.Second, + } + rctx, cancel := context.WithDeadline(ctx, time.Now().Add(30*time.Second)) + req, err := http.NewRequestWithContext(rctx, + http.MethodPost, url, bytes.NewReader(payload)) + defer cancel() + if err != nil { + return fmt.Errorf("http.NewRequestWithContext: %v: %e", + err, work.ErrDoNotReattempt) + } + + req.Header = headers + for key, values := range headers { + for _, value := range values { + req.Header.Add(key, value) + } + } + nonce, sig := crypto.SignWebhook(payload) + req.Header.Add("X-Payload-Nonce", nonce) + req.Header.Add("X-Payload-Signature", sig) + + resp, err := client.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + reader := io.LimitReader(resp.Body, 65536) // No more than 64 KiB + body, err := ioutil.ReadAll(reader) + if err != nil { + return fmt.Errorf("Error reading response body: %v: %e", + err, work.ErrDoNotReattempt) + } + + if err = database.WithTx(ctx, nil, func(tx *sql.Tx) error { + var sb strings.Builder + resp.Header.Write(&sb) + _, err := sq.Update(name+"_webhook_delivery"). + Set("response", string(body)). + Set("response_status", resp.StatusCode). + Set("response_headers", sb.String()). + RunWith(tx). + ExecContext(ctx) + return err + }); err != nil { + log.Printf("Warning: webhook delivered, but updating delivery record failed: %v", err) + return nil + } + + if resp.StatusCode == http.StatusBadGateway || + resp.StatusCode == http.StatusServiceUnavailable || + resp.StatusCode == http.StatusGatewayTimeout { + // Retry + return fmt.Errorf("Server returned status %d: %s", + resp.StatusCode, resp.Status) + } + + return nil +} diff --git a/webhooks/legacy_test.go b/webhooks/legacy_test.go new file mode 100644 index 0000000000000000000000000000000000000000..26cb27276840eb90c5f4fa3e9d083841c40ee580 --- /dev/null +++ b/webhooks/legacy_test.go @@ -0,0 +1,114 @@ +package webhooks + +import ( + "context" + "io/ioutil" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/DATA-DOG/go-sqlmock" + "github.com/stretchr/testify/assert" + "github.com/vaughan0/go-ini" + sq "github.com/Masterminds/squirrel" + + "git.sr.ht/~sircmpwn/core-go/crypto" + "git.sr.ht/~sircmpwn/core-go/database" +) + +func init() { + conf, err := ini.Load(strings.NewReader(` +[webhooks] +private-key=ebzsjPaN6E13ln/FeNWly1C92q6bVMVdOnDo1HPl5fc= + +[sr.ht] +network-key=tbuG-7Vh44vrDq1L_HKWkHnWrDOtJhEkPKPiauaLeuk= + +[test::api] +internal-ipnet=127.0.0.1/24,::1/64`)) + if err != nil { + panic(err) + } + crypto.InitCrypto(conf) +} + +func TestDelivery(t *testing.T) { + var called bool + 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") + + assert.NotEqual(t, "", r.Header.Get("X-Webhook-Delivery")) + assert.Equal(t, "profile:update", r.Header.Get("X-Webhook-Event")) + assert.Equal(t, "application/json", r.Header.Get("Content-Type")) + + b, err := ioutil.ReadAll(r.Body) + assert.Nil(t, err) + assert.Equal(t, `{"hello": "world"}`, string(b)) + + nonce := r.Header.Get("X-Payload-Nonce") + signature := r.Header.Get("X-Payload-Signature") + assert.True(t, crypto.VerifyWebhook(b, nonce, signature)) + + w.Write([]byte("Thanks!")) + })) + defer srv.Close() + + queue := NewLegacyQueue() + q := sq. + Select(). + From("user_webhook_subscription"). + Where(`user_id = ?`, 42) + queue.Schedule(q, "user", "profile:update", []byte(`{"hello": "world"}`)) + + db, mock, err := sqlmock.New() + if err != nil { + panic(err) + } + + // Lookup phase + mock.ExpectBegin() + mock.ExpectQuery(`SELECT .* FROM user_webhook_subscription`). + WillReturnRows(sqlmock.NewRows([]string{ + "id", "created", "url", "events", + }).AddRow( + 1337, time.Now().UTC(), + srv.URL + "/webhook", + "profile:update")). + WithArgs(42, sqlmock.AnyArg()) // Any => events LIKE %profile:update% + mock.ExpectCommit() + + // Schedule phase + mock.ExpectBegin() + mock.ExpectQuery(`INSERT INTO user_webhook_delivery`). + WillReturnRows(sqlmock.NewRows([]string{"id"}).AddRow(4096)) + mock.ExpectCommit() + + ctx := database.Context(context.Background(), db) + 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, sqlmock.AnyArg()) // Any => response headers + mock.ExpectCommit() + + ctx = database.Context(context.Background(), db) + queue.Queue.Dispatch(ctx) + + assert.Nil(t, mock.ExpectationsWereMet()) + assert.True(t, called) +}