diff --git a/spindle/webhook/webhook.go b/spindle/webhook/webhook.go new file mode 100644 index 000000000..fc438bca5 --- /dev/null +++ b/spindle/webhook/webhook.go @@ -0,0 +1,303 @@ +// Package webhook delivers repository event notifications to user-configured +// HTTP endpoints. It is the spindle-side successor to the appview webhook +// notifier: spindle owns webhook configuration and delivery, firing from the +// repo events it already ingests (push, pull_request:*, repository:renamed). +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/bluesky-social/indigo/atproto/syntax" + "github.com/google/uuid" + "tangled.org/core/api/tangled" + "tangled.org/core/hostutil" + "tangled.org/core/log" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" +) + +type Service struct { + db *db.DB + logger *slog.Logger + client *http.Client +} + +func New(database *db.DB, dev bool) *Service { + return &Service{ + db: database, + logger: log.New("webhook"), + // user-supplied webhook URLs are untrusted: block internal address + // ranges and don't follow redirects to guard against SSRF. + client: hostutil.SafeClient(dev, 30*time.Second), + } +} + +// FirePush delivers a push (git ref-update) event to subscribed webhooks. +func (s *Service) FirePush(ctx context.Context, repo *db.Repo, committerDid, ref, oldSha, newSha string) { + webhooks := s.activeWebhooksForEvent(repo.RepoDid, models.WebhookEventPush) + if len(webhooks) == 0 { + return + } + + pusher := committerDid + if pusher == "" { + pusher = repo.Owner.String() + } + payload := &models.WebhookPayload{ + Ref: ref, + Before: oldSha, + After: newSha, + Repository: buildRepository(repo), + Pusher: models.WebhookUser{Did: pusher}, + } + payloadBytes, err := json.Marshal(payload) + if err != nil { + s.logger.Error("failed to marshal push payload", "repo_did", repo.RepoDid, "err", err) + return + } + + userAgent := "Tangled-Hook/push" + if len(newSha) >= 7 { + userAgent = "Tangled-Hook/" + newSha[:7] + } + for _, webhook := range webhooks { + go s.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload.Repository.FullName, userAgent, payloadBytes) + } +} + +// FireRename delivers a repository:renamed event to subscribed webhooks. +func (s *Service) FireRename(ctx context.Context, actor syntax.DID, repo *db.Repo, oldName, newName string) { + webhooks := s.activeWebhooksForEvent(repo.RepoDid, models.WebhookEventRepoRenamed) + if len(webhooks) == 0 { + return + } + + payload := &models.WebhookRenamePayload{ + OldName: oldName, + NewName: newName, + Repository: buildRepository(repo), + Sender: models.WebhookUser{Did: actor.String()}, + } + payloadBytes, err := json.Marshal(payload) + if err != nil { + s.logger.Error("failed to marshal rename payload", "repo_did", repo.RepoDid, "err", err) + return + } + + for _, webhook := range webhooks { + go s.sendWebhook(ctx, webhook, string(models.WebhookEventRepoRenamed), payload.Repository.FullName, "Tangled-Hook/rename", payloadBytes) + } +} + +// FirePullRequest delivers a pull_request:* event to subscribed webhooks. repo +// is the target repository; record is the pull record; pullAuthor authored the +// pull; sender is the actor that triggered this transition (equal to pullAuthor +// for created/resubmitted). +func (s *Service) FirePullRequest(ctx context.Context, event models.WebhookEvent, action string, repo *db.Repo, record *tangled.RepoPull, pullAuthor, sender syntax.DID) { + // pull request events may originate from firehose handlers whose context is + // canceled promptly; detach so in-flight deliveries are not cut short + ctx = context.WithoutCancel(ctx) + + webhooks := s.activeWebhooksForEvent(repo.RepoDid, event) + if len(webhooks) == 0 { + return + } + + pr := models.WebhookPullRequest{ + Title: record.Title, + State: pullStateFromAction(action), + RoundNumber: len(record.Rounds), + Owner: models.WebhookUser{Did: pullAuthor.String()}, + CreatedAt: record.CreatedAt, + } + if record.Body != nil { + pr.Body = *record.Body + } + if record.Target != nil { + pr.TargetBranch = record.Target.Branch + } + if record.Source != nil { + source := &models.WebhookPullRequestSource{Branch: record.Source.Branch} + if record.Source.Repo != nil { + source.Repo = *record.Source.Repo + } + pr.Source = source + } + + payload := &models.WebhookPullRequestPayload{ + Action: action, + PullRequest: pr, + Repository: buildRepository(repo), + Sender: models.WebhookUser{Did: sender.String()}, + } + payloadBytes, err := json.Marshal(payload) + if err != nil { + s.logger.Error("failed to marshal pull request payload", "repo_did", repo.RepoDid, "err", err) + return + } + + for _, webhook := range webhooks { + go s.sendWebhook(ctx, webhook, string(event), payload.Repository.FullName, "Tangled-Hook/pull_request", payloadBytes) + } +} + +// Redeliver re-sends a stored delivery via the live send-and-record path, +// signing with the webhook's current secret. +func (s *Service) Redeliver(ctx context.Context, webhook models.Webhook, prev models.WebhookDelivery) { + // recover the repo full name (X-Tangled-Repo header) from the stored payload + var meta struct { + Repository struct { + FullName string `json:"full_name"` + } `json:"repository"` + } + _ = json.Unmarshal([]byte(prev.RequestBody), &meta) + + s.sendWebhook(ctx, webhook, prev.Event, meta.Repository.FullName, "Tangled-Hook/retry", []byte(prev.RequestBody)) +} + +// pullStateFromAction derives the pull request state carried in the payload +// from the lifecycle action, since the pull record does not store state. +func pullStateFromAction(action string) string { + switch action { + case "merged": + return "merged" + case "closed": + return "closed" + default: // created, resubmitted, reopened + return "open" + } +} + +func (s *Service) activeWebhooksForEvent(repoDid syntax.DID, event models.WebhookEvent) []models.Webhook { + webhooks, err := s.db.GetActiveWebhooksForRepo(repoDid) + if err != nil { + s.logger.Error("failed to get webhooks for repo", "repo_did", repoDid, "err", err) + return nil + } + var matching []models.Webhook + for _, webhook := range webhooks { + if webhook.HasEvent(event) { + matching = append(matching, webhook) + } + } + return matching +} + +func buildRepository(repo *db.Repo) models.WebhookRepository { + owner := repo.Owner.String() + rkey := repo.Rkey.String() + return models.WebhookRepository{ + Name: repo.Name, + FullName: fmt.Sprintf("%s/%s", owner, rkey), + HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, owner, rkey), + CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, owner, rkey), + SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, owner, rkey), + CreatedAt: repo.CreatedAt, + UpdatedAt: repo.CreatedAt, + Owner: models.WebhookUser{Did: owner}, + } +} + +func (s *Service) sendWebhook(ctx context.Context, webhook models.Webhook, event, repoFullName, userAgent string, payloadBytes []byte) { + deliveryId := uuid.New().String() + + var signature string + if webhook.Secret != "" { + signature = "sha256=" + s.computeSignature(payloadBytes, webhook.Secret) + } + + 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) { + s.logger.Info("retrying webhook delivery", "webhook_id", webhook.Id, "attempt", n+1, "err", err) + }), + retry.Context(ctx), + } + + var resp *http.Response + err := retry.Do(func() error { + // build a fresh request each attempt so the body reader is not + // exhausted after the first try + req, err := http.NewRequestWithContext(ctx, "POST", webhook.Url, bytes.NewReader(payloadBytes)) + if err != nil { + return retry.Unrecoverable(err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", userAgent) + 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", repoFullName) + if signature != "" { + req.Header.Set("X-Tangled-Signature-256", signature) + } + + r, err := s.client.Do(req) + if err != nil { + return err + } + if r.StatusCode >= 500 { + r.Body.Close() + return fmt.Errorf("server error: %d", r.StatusCode) + } + resp = r + return nil + }, retryOpts...) + + if err != nil { + s.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 { + s.logger.Warn("failed to read webhook response body", "webhook_id", webhook.Id, "err", err) + } else { + delivery.ResponseBody = string(bodyBytes) + } + + if !delivery.Success { + s.logger.Warn("webhook delivery failed", "webhook_id", webhook.Id, "status", resp.StatusCode, "url", webhook.Url) + } else { + s.logger.Info("webhook delivered successfully", "webhook_id", webhook.Id, "url", webhook.Url, "delivery_id", deliveryId) + } + } + + if err := s.db.AddWebhookDelivery(delivery); err != nil { + s.logger.Error("failed to record webhook delivery", "webhook_id", webhook.Id, "err", err) + } +} + +func (s *Service) computeSignature(payload []byte, secret string) string { + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + return hex.EncodeToString(mac.Sum(nil)) +}