// 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" "errors" "fmt" "io" "log/slog" "net/http" "sync" "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" ) const ( deliveryWorkers = 8 deliveryQueueSize = 256 deliveryRetention = 50 ) var errQueueSaturated = errors.New("delivery queue saturated") type Service struct { db *db.DB logger *slog.Logger client *http.Client jobs chan deliveryJob wg sync.WaitGroup closeOnce sync.Once } func New(database *db.DB, dev bool) *Service { s := &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), jobs: make(chan deliveryJob, deliveryQueueSize), } for range deliveryWorkers { s.wg.Add(1) go s.worker() } return s } // 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 { s.dispatch(deliveryJob{ctx: ctx, webhook: webhook, event: string(models.WebhookEventPush), repoFullName: payload.Repository.FullName, userAgent: userAgent, payload: 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 { s.dispatch(deliveryJob{ctx: ctx, webhook: webhook, event: string(models.WebhookEventRepoRenamed), repoFullName: payload.Repository.FullName, userAgent: "Tangled-Hook/rename", payload: 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 { s.dispatch(deliveryJob{ctx: ctx, webhook: webhook, event: string(event), repoFullName: payload.Repository.FullName, userAgent: "Tangled-Hook/pull_request", payload: 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.dispatch(deliveryJob{ctx: ctx, webhook: webhook, event: prev.Event, repoFullName: meta.Repository.FullName, userAgent: "Tangled-Hook/retry", payload: []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}, } } type deliveryJob struct { ctx context.Context webhook models.Webhook event string repoFullName string userAgent string payload []byte } func (s *Service) dispatch(job deliveryJob) { select { case s.jobs <- job: default: s.logger.Warn("webhook delivery queue saturated", "webhook_id", job.webhook.Id, "url", job.webhook.Url) s.record(&models.WebhookDelivery{ WebhookId: job.webhook.Id, Event: job.event, DeliveryId: uuid.New().String(), Url: job.webhook.Url, RequestBody: string(job.payload), ResponseBody: errQueueSaturated.Error(), }) } } func (s *Service) worker() { defer s.wg.Done() for job := range s.jobs { s.sendWebhook(job) } } func (s *Service) Close(ctx context.Context) { s.closeOnce.Do(func() { close(s.jobs) }) done := make(chan struct{}) go func() { s.wg.Wait() close(done) }() select { case <-done: case <-ctx.Done(): s.logger.Warn("webhook deliveries still in flight at shutdown") } } func (s *Service) sendWebhook(job deliveryJob) { ctx := job.ctx var signature string if job.webhook.Secret != "" { signature = "sha256=" + s.computeSignature(job.payload, job.webhook.Secret) } deliveryId := uuid.New().String() 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", job.webhook.Id, "attempt", n+1, "err", err) }), retry.Context(ctx), } var ( resp *http.Response lastStatus int lastResponse string ) 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", job.webhook.Url, bytes.NewReader(job.payload)) if err != nil { return retry.Unrecoverable(err) } req.Header.Set("Content-Type", "application/json") req.Header.Set("User-Agent", job.userAgent) req.Header.Set("X-Tangled-Event", job.event) req.Header.Set("X-Tangled-Hook-ID", fmt.Sprintf("%d", job.webhook.Id)) req.Header.Set("X-Tangled-Delivery", deliveryId) req.Header.Set("X-Tangled-Repo", job.repoFullName) if signature != "" { req.Header.Set("X-Tangled-Signature-256", signature) } r, err := s.client.Do(req) if err != nil { return err } lastStatus = r.StatusCode body, readErr := io.ReadAll(io.LimitReader(r.Body, 10*1024)) r.Body.Close() if readErr != nil { s.logger.Warn("failed to read webhook response body", "webhook_id", job.webhook.Id, "err", readErr) } lastResponse = string(body) if r.StatusCode >= 500 { return fmt.Errorf("server error: %d", r.StatusCode) } resp = r return nil }, retryOpts...) delivery := &models.WebhookDelivery{ WebhookId: job.webhook.Id, Event: job.event, DeliveryId: deliveryId, Url: job.webhook.Url, RequestBody: string(job.payload), ResponseCode: lastStatus, ResponseBody: lastResponse, } switch { case err != nil: s.logger.Error("webhook request failed after retries", "webhook_id", job.webhook.Id, "err", err) delivery.ResponseBody = err.Error() case resp != nil && resp.StatusCode >= 200 && resp.StatusCode < 300: delivery.Success = true s.logger.Info("webhook delivered successfully", "webhook_id", job.webhook.Id, "url", job.webhook.Url, "delivery_id", deliveryId) default: s.logger.Warn("webhook delivery failed", "webhook_id", job.webhook.Id, "status", lastStatus, "url", job.webhook.Url) } s.record(delivery) } func (s *Service) record(delivery *models.WebhookDelivery) { if err := s.db.AddWebhookDelivery(delivery); err != nil { s.logger.Error("failed to record webhook delivery", "webhook_id", delivery.WebhookId, "err", err) return } if err := s.db.PruneWebhookDeliveries(delivery.WebhookId, deliveryRetention); err != nil { s.logger.Error("failed to prune webhook deliveries", "webhook_id", delivery.WebhookId, "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)) }