Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378// 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))}