From 3048b68ecb8934e5625825443965610d1266d519 Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Thu, 02 Apr 2026 08:37:08 +0000 Subject: [PATCH] appview/notify: move logging and webhook notifiers to their own packages Signed-off-by: Anirudh Oppiliappan --- appview/notify/logging_notifier.go | 125 ----------------------------------------------------------------------------------------------------------------------------- appview/notify/webhook_notifier.go | 241 ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- appview/state/state.go | 7 ++++--- appview/notify/logging/notifier.go | 122 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ appview/notify/webhook/notifier.go | 224 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 5 file(s) changed, 350 insertion(s)(+), 369 deletion(s)(-) diff --git a/appview/notify/logging_notifier.go b/appview/notify/logging_notifier.go deleted file mode 100644 --- a/appview/notify/logging_notifier.go +++ /dev/null @@ -1,125 +0,0 @@ -package notify - -import ( - "context" - "log/slog" - - "tangled.org/core/appview/models" - tlog "tangled.org/core/log" - - "github.com/bluesky-social/indigo/atproto/syntax" -) - -type loggingNotifier struct { - inner Notifier - logger *slog.Logger -} - -func NewLoggingNotifier(inner Notifier, logger *slog.Logger) Notifier { - return &loggingNotifier{ - inner, - logger, - } -} - -var _ Notifier = &loggingNotifier{} - -func (l *loggingNotifier) NewRepo(ctx context.Context, repo *models.Repo) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewRepo")) - l.inner.NewRepo(ctx, repo) -} - -func (l *loggingNotifier) NewStar(ctx context.Context, star *models.Star) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewStar")) - l.inner.NewStar(ctx, star) -} - -func (l *loggingNotifier) DeleteStar(ctx context.Context, star *models.Star) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteStar")) - l.inner.DeleteStar(ctx, star) -} - -func (l *loggingNotifier) NewIssue(ctx context.Context, issue *models.Issue, mentions []syntax.DID) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssue")) - l.inner.NewIssue(ctx, issue, mentions) -} - -func (l *loggingNotifier) NewIssueComment(ctx context.Context, comment *models.IssueComment, mentions []syntax.DID) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueComment")) - l.inner.NewIssueComment(ctx, comment, mentions) -} - -func (l *loggingNotifier) NewIssueState(ctx context.Context, actor syntax.DID, issue *models.Issue) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueState")) - l.inner.NewIssueState(ctx, actor, issue) -} - -func (l *loggingNotifier) DeleteIssue(ctx context.Context, issue *models.Issue) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteIssue")) - l.inner.DeleteIssue(ctx, issue) -} - -func (l *loggingNotifier) NewIssueLabelOp(ctx context.Context, issue *models.Issue) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueLabelOp")) - l.inner.NewIssueLabelOp(ctx, issue) -} - -func (l *loggingNotifier) NewPullLabelOp(ctx context.Context, pull *models.Pull) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullLabelOp")) - l.inner.NewPullLabelOp(ctx, pull) -} - -func (l *loggingNotifier) NewFollow(ctx context.Context, follow *models.Follow) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewFollow")) - l.inner.NewFollow(ctx, follow) -} - -func (l *loggingNotifier) DeleteFollow(ctx context.Context, follow *models.Follow) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteFollow")) - l.inner.DeleteFollow(ctx, follow) -} - -func (l *loggingNotifier) NewPull(ctx context.Context, pull *models.Pull) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPull")) - l.inner.NewPull(ctx, pull) -} - -func (l *loggingNotifier) NewPullComment(ctx context.Context, comment *models.PullComment, mentions []syntax.DID) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullComment")) - l.inner.NewPullComment(ctx, comment, mentions) -} - -func (l *loggingNotifier) NewPullState(ctx context.Context, actor syntax.DID, pull *models.Pull) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullState")) - l.inner.NewPullState(ctx, actor, pull) -} - -func (l *loggingNotifier) UpdateProfile(ctx context.Context, profile *models.Profile) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "UpdateProfile")) - l.inner.UpdateProfile(ctx, profile) -} - -func (l *loggingNotifier) NewString(ctx context.Context, s *models.String) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewString")) - l.inner.NewString(ctx, s) -} - -func (l *loggingNotifier) EditString(ctx context.Context, s *models.String) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "EditString")) - l.inner.EditString(ctx, s) -} - -func (l *loggingNotifier) DeleteString(ctx context.Context, did, rkey string) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteString")) - l.inner.DeleteString(ctx, did, rkey) -} - -func (l *loggingNotifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "Push")) - l.inner.Push(ctx, repo, ref, oldSha, newSha, committerDid) -} - -func (l *loggingNotifier) Clone(ctx context.Context, repo *models.Repo) { - ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "Clone")) - l.inner.Clone(ctx, repo) -} diff --git a/appview/notify/webhook_notifier.go b/appview/notify/webhook_notifier.go deleted file mode 100644 --- a/appview/notify/webhook_notifier.go +++ /dev/null @@ -1,241 +0,0 @@ -package notify - -import ( - "bytes" - "context" - "crypto/hmac" - "crypto/sha256" - "encoding/hex" - "encoding/json" - "fmt" - "io" - "log/slog" - "net/http" - "time" - - "github.com/avast/retry-go/v4" - "github.com/google/uuid" - "tangled.org/core/appview/db" - "tangled.org/core/appview/models" - "tangled.org/core/log" -) - -type WebhookNotifier struct { - BaseNotifier - db *db.DB - logger *slog.Logger - client *http.Client -} - -func NewWebhookNotifier(database *db.DB) *WebhookNotifier { - return &WebhookNotifier{ - db: database, - logger: log.New("webhook-notifier"), - client: &http.Client{ - Timeout: 30 * time.Second, - }, - } -} - -// Push implements the Notifier interface for git push events -func (w *WebhookNotifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { - webhooks, err := db.GetActiveWebhooksForRepo(w.db, repo.RepoAt()) - if err != nil { - w.logger.Error("failed to get webhooks for repo", "repo", repo.RepoAt(), "err", err) - return - } - - // check if any webhooks are subscribed to push events - var pushWebhooks []models.Webhook - for _, webhook := range webhooks { - if webhook.HasEvent(models.WebhookEventPush) { - pushWebhooks = append(pushWebhooks, webhook) - } - } - - if len(pushWebhooks) == 0 { - return - } - - payload, err := w.buildPushPayload(repo, ref, oldSha, newSha, committerDid) - if err != nil { - w.logger.Error("failed to build push payload", "repo", repo.RepoAt(), "err", err) - return - } - - // Send webhooks - for _, webhook := range pushWebhooks { - go w.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload) - } -} - -func (w *WebhookNotifier) Clone(ctx context.Context, repo *models.Repo) {} - -// buildPushPayload creates the webhook payload -func (w *WebhookNotifier) buildPushPayload(repo *models.Repo, ref, oldSha, newSha, committerDid string) (*models.WebhookPayload, error) { - owner := repo.Did - - pusher := committerDid - if committerDid == "" { - pusher = owner - } - - // Build repository object - repository := models.WebhookRepository{ - Name: repo.Name, - FullName: fmt.Sprintf("%s/%s", repo.Did, repo.Name), - Description: repo.Description, - Fork: repo.Source != "", - HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), - CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), - SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, repo.Did, repo.Name), - CreatedAt: repo.Created.Format(time.RFC3339), - UpdatedAt: repo.Created.Format(time.RFC3339), - Owner: models.WebhookUser{ - Did: owner, - }, - } - - // Add optional fields - if repo.Website != "" { - repository.Website = repo.Website - } - if repo.RepoStats != nil { - repository.StarsCount = repo.RepoStats.StarCount - repository.OpenIssues = repo.RepoStats.IssueCount.Open - } - - // Build payload - payload := &models.WebhookPayload{ - Ref: ref, - Before: oldSha, - After: newSha, - Repository: repository, - Pusher: models.WebhookUser{ - Did: pusher, - }, - } - - return payload, nil -} - -// sendWebhook sends the webhook http request -func (w *WebhookNotifier) sendWebhook(ctx context.Context, webhook models.Webhook, event string, payload *models.WebhookPayload) { - deliveryId := uuid.New().String() - - payloadBytes, err := json.Marshal(payload) - if err != nil { - w.logger.Error("failed to marshal webhook payload", "webhook_id", webhook.Id, "err", err) - return - } - - req, err := http.NewRequestWithContext(ctx, "POST", webhook.Url, bytes.NewReader(payloadBytes)) - if err != nil { - w.logger.Error("failed to create webhook request", "webhook_id", webhook.Id, "err", err) - return - } - - shortSha := payload.After[:7] - - req.Header.Set("Content-Type", "application/json") - req.Header.Set("User-Agent", "Tangled-Hook/"+shortSha) - req.Header.Set("X-Tangled-Event", event) - req.Header.Set("X-Tangled-Hook-ID", fmt.Sprintf("%d", webhook.Id)) - req.Header.Set("X-Tangled-Delivery", deliveryId) - req.Header.Set("X-Tangled-Repo", payload.Repository.FullName) - - if webhook.Secret != "" { - signature := w.computeSignature(payloadBytes, webhook.Secret) - req.Header.Set("X-Tangled-Signature-256", "sha256="+signature) - } - - delivery := &models.WebhookDelivery{ - WebhookId: webhook.Id, - Event: event, - DeliveryId: deliveryId, - Url: webhook.Url, - RequestBody: string(payloadBytes), - } - - // retry webhook delivery with exponential backoff - retryOpts := []retry.Option{ - retry.Attempts(3), - retry.Delay(1 * time.Second), - retry.MaxDelay(10 * time.Second), - retry.DelayType(retry.BackOffDelay), - retry.LastErrorOnly(true), - retry.OnRetry(func(n uint, err error) { - w.logger.Info("retrying webhook delivery", - "webhook_id", webhook.Id, - "attempt", n+1, - "err", err) - }), - retry.Context(ctx), - retry.RetryIf(func(err error) bool { - // only retry on network errors or 5xx responses - if err != nil { - return true - } - return false - }), - } - - var resp *http.Response - err = retry.Do(func() error { - var err error - resp, err = w.client.Do(req) - if err != nil { - return err - } - - // retry on 5xx server errors - if resp.StatusCode >= 500 { - defer resp.Body.Close() - return fmt.Errorf("server error: %d", resp.StatusCode) - } - - return nil - }, retryOpts...) - - if err != nil { - w.logger.Error("webhook request failed after retries", "webhook_id", webhook.Id, "err", err) - delivery.Success = false - delivery.ResponseBody = err.Error() - } else { - defer resp.Body.Close() - - delivery.ResponseCode = resp.StatusCode - delivery.Success = resp.StatusCode >= 200 && resp.StatusCode < 300 - - // Read response body (limit to 10KB) - bodyBytes, err := io.ReadAll(io.LimitReader(resp.Body, 10*1024)) - if err != nil { - w.logger.Warn("failed to read webhook response body", "webhook_id", webhook.Id, "err", err) - } else { - delivery.ResponseBody = string(bodyBytes) - } - - if !delivery.Success { - w.logger.Warn("webhook delivery failed", - "webhook_id", webhook.Id, - "status", resp.StatusCode, - "url", webhook.Url) - } else { - w.logger.Info("webhook delivered successfully", - "webhook_id", webhook.Id, - "url", webhook.Url, - "delivery_id", deliveryId) - } - } - - if err := db.AddWebhookDelivery(w.db, delivery); err != nil { - w.logger.Error("failed to record webhook delivery", "webhook_id", webhook.Id, "err", err) - } -} - -// computeSignature computes HMAC-SHA256 signature for the payload -func (w *WebhookNotifier) computeSignature(payload []byte, secret string) string { - mac := hmac.New(sha256.New, []byte(secret)) - mac.Write(payload) - return hex.EncodeToString(mac.Sum(nil)) -} diff --git a/appview/state/state.go b/appview/state/state.go --- a/appview/state/state.go +++ b/appview/state/state.go @@ -21,7 +21,9 @@ "tangled.org/core/appview/models" "tangled.org/core/appview/notify" dbnotify "tangled.org/core/appview/notify/db" + lognotify "tangled.org/core/appview/notify/logging" phnotify "tangled.org/core/appview/notify/posthog" + whnotify "tangled.org/core/appview/notify/webhook" "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" "tangled.org/core/appview/reporesolver" @@ -168,11 +170,10 @@ } notifiers = append(notifiers, indexer) - // Add webhook notifier - notifiers = append(notifiers, notify.NewWebhookNotifier(d)) + notifiers = append(notifiers, whnotify.NewNotifier(d)) notifier := notify.NewMergedNotifier(notifiers) - notifier = notify.NewLoggingNotifier(notifier, tlog.SubLogger(logger, "notify")) + notifier = lognotify.NewLoggingNotifier(notifier, tlog.SubLogger(logger, "notify")) var cfClient *cloudflare.Client if config.Cloudflare.ApiToken != "" { diff --git a/appview/notify/logging/notifier.go b/appview/notify/logging/notifier.go new file mode 100644 --- /dev/null +++ b/appview/notify/logging/notifier.go @@ -0,0 +1,122 @@ +package logging + +import ( + "context" + "log/slog" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/appview/models" + "tangled.org/core/appview/notify" + tlog "tangled.org/core/log" +) + +type loggingNotifier struct { + inner notify.Notifier + logger *slog.Logger +} + +func NewLoggingNotifier(inner notify.Notifier, logger *slog.Logger) notify.Notifier { + return &loggingNotifier{inner, logger} +} + +var _ notify.Notifier = &loggingNotifier{} + +func (l *loggingNotifier) NewRepo(ctx context.Context, repo *models.Repo) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewRepo")) + l.inner.NewRepo(ctx, repo) +} + +func (l *loggingNotifier) NewStar(ctx context.Context, star *models.Star) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewStar")) + l.inner.NewStar(ctx, star) +} + +func (l *loggingNotifier) DeleteStar(ctx context.Context, star *models.Star) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteStar")) + l.inner.DeleteStar(ctx, star) +} + +func (l *loggingNotifier) NewIssue(ctx context.Context, issue *models.Issue, mentions []syntax.DID) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssue")) + l.inner.NewIssue(ctx, issue, mentions) +} + +func (l *loggingNotifier) NewIssueComment(ctx context.Context, comment *models.IssueComment, mentions []syntax.DID) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueComment")) + l.inner.NewIssueComment(ctx, comment, mentions) +} + +func (l *loggingNotifier) NewIssueState(ctx context.Context, actor syntax.DID, issue *models.Issue) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueState")) + l.inner.NewIssueState(ctx, actor, issue) +} + +func (l *loggingNotifier) DeleteIssue(ctx context.Context, issue *models.Issue) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteIssue")) + l.inner.DeleteIssue(ctx, issue) +} + +func (l *loggingNotifier) NewIssueLabelOp(ctx context.Context, issue *models.Issue) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewIssueLabelOp")) + l.inner.NewIssueLabelOp(ctx, issue) +} + +func (l *loggingNotifier) NewPullLabelOp(ctx context.Context, pull *models.Pull) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullLabelOp")) + l.inner.NewPullLabelOp(ctx, pull) +} + +func (l *loggingNotifier) NewFollow(ctx context.Context, follow *models.Follow) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewFollow")) + l.inner.NewFollow(ctx, follow) +} + +func (l *loggingNotifier) DeleteFollow(ctx context.Context, follow *models.Follow) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteFollow")) + l.inner.DeleteFollow(ctx, follow) +} + +func (l *loggingNotifier) NewPull(ctx context.Context, pull *models.Pull) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPull")) + l.inner.NewPull(ctx, pull) +} + +func (l *loggingNotifier) NewPullComment(ctx context.Context, comment *models.PullComment, mentions []syntax.DID) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullComment")) + l.inner.NewPullComment(ctx, comment, mentions) +} + +func (l *loggingNotifier) NewPullState(ctx context.Context, actor syntax.DID, pull *models.Pull) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewPullState")) + l.inner.NewPullState(ctx, actor, pull) +} + +func (l *loggingNotifier) UpdateProfile(ctx context.Context, profile *models.Profile) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "UpdateProfile")) + l.inner.UpdateProfile(ctx, profile) +} + +func (l *loggingNotifier) NewString(ctx context.Context, s *models.String) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "NewString")) + l.inner.NewString(ctx, s) +} + +func (l *loggingNotifier) EditString(ctx context.Context, s *models.String) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "EditString")) + l.inner.EditString(ctx, s) +} + +func (l *loggingNotifier) DeleteString(ctx context.Context, did, rkey string) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "DeleteString")) + l.inner.DeleteString(ctx, did, rkey) +} + +func (l *loggingNotifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "Push")) + l.inner.Push(ctx, repo, ref, oldSha, newSha, committerDid) +} + +func (l *loggingNotifier) Clone(ctx context.Context, repo *models.Repo) { + ctx = tlog.IntoContext(ctx, tlog.SubLogger(l.logger, "Clone")) + l.inner.Clone(ctx, repo) +} diff --git a/appview/notify/webhook/notifier.go b/appview/notify/webhook/notifier.go new file mode 100644 --- /dev/null +++ b/appview/notify/webhook/notifier.go @@ -0,0 +1,224 @@ +package webhook + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "log/slog" + "net/http" + "time" + + "github.com/avast/retry-go/v4" + "github.com/google/uuid" + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/appview/notify" + "tangled.org/core/log" +) + +type Notifier struct { + notify.BaseNotifier + db *db.DB + logger *slog.Logger + client *http.Client +} + +func NewNotifier(database *db.DB) *Notifier { + return &Notifier{ + db: database, + logger: log.New("webhook-notifier"), + client: &http.Client{ + Timeout: 30 * time.Second, + }, + } +} + +var _ notify.Notifier = &Notifier{} + +func (w *Notifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { + webhooks, err := db.GetActiveWebhooksForRepo(w.db, repo.RepoAt()) + if err != nil { + w.logger.Error("failed to get webhooks for repo", "repo", repo.RepoAt(), "err", err) + return + } + + var pushWebhooks []models.Webhook + for _, webhook := range webhooks { + if webhook.HasEvent(models.WebhookEventPush) { + pushWebhooks = append(pushWebhooks, webhook) + } + } + + if len(pushWebhooks) == 0 { + return + } + + payload, err := w.buildPushPayload(repo, ref, oldSha, newSha, committerDid) + if err != nil { + w.logger.Error("failed to build push payload", "repo", repo.RepoAt(), "err", err) + return + } + + for _, webhook := range pushWebhooks { + go w.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload) + } +} + +func (w *Notifier) buildPushPayload(repo *models.Repo, ref, oldSha, newSha, committerDid string) (*models.WebhookPayload, error) { + owner := repo.Did + + pusher := committerDid + if committerDid == "" { + pusher = owner + } + + repository := models.WebhookRepository{ + Name: repo.Name, + FullName: fmt.Sprintf("%s/%s", repo.Did, repo.Name), + Description: repo.Description, + Fork: repo.Source != "", + HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), + CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), + SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, repo.Did, repo.Name), + CreatedAt: repo.Created.Format(time.RFC3339), + UpdatedAt: repo.Created.Format(time.RFC3339), + Owner: models.WebhookUser{ + Did: owner, + }, + } + + if repo.Website != "" { + repository.Website = repo.Website + } + if repo.RepoStats != nil { + repository.StarsCount = repo.RepoStats.StarCount + repository.OpenIssues = repo.RepoStats.IssueCount.Open + } + + payload := &models.WebhookPayload{ + Ref: ref, + Before: oldSha, + After: newSha, + Repository: repository, + Pusher: models.WebhookUser{ + Did: pusher, + }, + } + + return payload, nil +} + +func (w *Notifier) sendWebhook(ctx context.Context, webhook models.Webhook, event string, payload *models.WebhookPayload) { + deliveryId := uuid.New().String() + + payloadBytes, err := json.Marshal(payload) + if err != nil { + w.logger.Error("failed to marshal webhook payload", "webhook_id", webhook.Id, "err", err) + return + } + + req, err := http.NewRequestWithContext(ctx, "POST", webhook.Url, bytes.NewReader(payloadBytes)) + if err != nil { + w.logger.Error("failed to create webhook request", "webhook_id", webhook.Id, "err", err) + return + } + + shortSha := payload.After[:7] + + req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", "Tangled-Hook/"+shortSha) + req.Header.Set("X-Tangled-Event", event) + req.Header.Set("X-Tangled-Hook-ID", fmt.Sprintf("%d", webhook.Id)) + req.Header.Set("X-Tangled-Delivery", deliveryId) + req.Header.Set("X-Tangled-Repo", payload.Repository.FullName) + + if webhook.Secret != "" { + signature := w.computeSignature(payloadBytes, webhook.Secret) + req.Header.Set("X-Tangled-Signature-256", "sha256="+signature) + } + + delivery := &models.WebhookDelivery{ + WebhookId: webhook.Id, + Event: event, + DeliveryId: deliveryId, + Url: webhook.Url, + RequestBody: string(payloadBytes), + } + + retryOpts := []retry.Option{ + retry.Attempts(3), + retry.Delay(1 * time.Second), + retry.MaxDelay(10 * time.Second), + retry.DelayType(retry.BackOffDelay), + retry.LastErrorOnly(true), + retry.OnRetry(func(n uint, err error) { + w.logger.Info("retrying webhook delivery", + "webhook_id", webhook.Id, + "attempt", n+1, + "err", err) + }), + retry.Context(ctx), + retry.RetryIf(func(err error) bool { + return err != nil + }), + } + + var resp *http.Response + err = retry.Do(func() error { + var err error + resp, err = w.client.Do(req) + if err != nil { + return err + } + if resp.StatusCode >= 500 { + defer resp.Body.Close() + return fmt.Errorf("server error: %d", resp.StatusCode) + } + return nil + }, retryOpts...) + + if err != nil { + w.logger.Error("webhook request failed after retries", "webhook_id", webhook.Id, "err", err) + delivery.Success = false + delivery.ResponseBody = err.Error() + } else { + defer resp.Body.Close() + + delivery.ResponseCode = resp.StatusCode + delivery.Success = resp.StatusCode >= 200 && resp.StatusCode < 300 + + bodyBytes, err := io.ReadAll(io.LimitReader(resp.Body, 10*1024)) + if err != nil { + w.logger.Warn("failed to read webhook response body", "webhook_id", webhook.Id, "err", err) + } else { + delivery.ResponseBody = string(bodyBytes) + } + + if !delivery.Success { + w.logger.Warn("webhook delivery failed", + "webhook_id", webhook.Id, + "status", resp.StatusCode, + "url", webhook.Url) + } else { + w.logger.Info("webhook delivered successfully", + "webhook_id", webhook.Id, + "url", webhook.Url, + "delivery_id", deliveryId) + } + } + + if err := db.AddWebhookDelivery(w.db, delivery); err != nil { + w.logger.Error("failed to record webhook delivery", "webhook_id", webhook.Id, "err", err) + } +} + +func (w *Notifier) computeSignature(payload []byte, secret string) string { + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + return hex.EncodeToString(mac.Sum(nil)) +} -- tangled.sh