This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / notify / webhook_notifier.go
6.5 kB 241 lines
1package notify 2 3import ( 4 "bytes" 5 "context" 6 "crypto/hmac" 7 "crypto/sha256" 8 "encoding/hex" 9 "encoding/json" 10 "fmt" 11 "io" 12 "log/slog" 13 "net/http" 14 "time" 15 16 "github.com/avast/retry-go/v4" 17 "github.com/google/uuid" 18 "tangled.org/core/appview/db" 19 "tangled.org/core/appview/models" 20 "tangled.org/core/log" 21) 22 23type WebhookNotifier struct { 24 BaseNotifier 25 db *db.DB 26 logger *slog.Logger 27 client *http.Client 28} 29 30func NewWebhookNotifier(database *db.DB) *WebhookNotifier { 31 return &WebhookNotifier{ 32 db: database, 33 logger: log.New("webhook-notifier"), 34 client: &http.Client{ 35 Timeout: 30 * time.Second, 36 }, 37 } 38} 39 40// Push implements the Notifier interface for git push events 41func (w *WebhookNotifier) Push(ctx context.Context, repo *models.Repo, ref, oldSha, newSha, committerDid string) { 42 webhooks, err := db.GetActiveWebhooksForRepo(w.db, repo.RepoAt()) 43 if err != nil { 44 w.logger.Error("failed to get webhooks for repo", "repo", repo.RepoAt(), "err", err) 45 return 46 } 47 48 // check if any webhooks are subscribed to push events 49 var pushWebhooks []models.Webhook 50 for _, webhook := range webhooks { 51 if webhook.HasEvent(models.WebhookEventPush) { 52 pushWebhooks = append(pushWebhooks, webhook) 53 } 54 } 55 56 if len(pushWebhooks) == 0 { 57 return 58 } 59 60 payload, err := w.buildPushPayload(repo, ref, oldSha, newSha, committerDid) 61 if err != nil { 62 w.logger.Error("failed to build push payload", "repo", repo.RepoAt(), "err", err) 63 return 64 } 65 66 // Send webhooks 67 for _, webhook := range pushWebhooks { 68 go w.sendWebhook(ctx, webhook, string(models.WebhookEventPush), payload) 69 } 70} 71 72func (w *WebhookNotifier) Clone(ctx context.Context, repo *models.Repo) {} 73 74// buildPushPayload creates the webhook payload 75func (w *WebhookNotifier) buildPushPayload(repo *models.Repo, ref, oldSha, newSha, committerDid string) (*models.WebhookPayload, error) { 76 owner := repo.Did 77 78 pusher := committerDid 79 if committerDid == "" { 80 pusher = owner 81 } 82 83 // Build repository object 84 repository := models.WebhookRepository{ 85 Name: repo.Name, 86 FullName: fmt.Sprintf("%s/%s", repo.Did, repo.Name), 87 Description: repo.Description, 88 Fork: repo.Source != "", 89 HtmlUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), 90 CloneUrl: fmt.Sprintf("https://%s/%s/%s", repo.Knot, repo.Did, repo.Name), 91 SshUrl: fmt.Sprintf("ssh://git@%s/%s/%s", repo.Knot, repo.Did, repo.Name), 92 CreatedAt: repo.Created.Format(time.RFC3339), 93 UpdatedAt: repo.Created.Format(time.RFC3339), 94 Owner: models.WebhookUser{ 95 Did: owner, 96 }, 97 } 98 99 // Add optional fields 100 if repo.Website != "" { 101 repository.Website = repo.Website 102 } 103 if repo.RepoStats != nil { 104 repository.StarsCount = repo.RepoStats.StarCount 105 repository.OpenIssues = repo.RepoStats.IssueCount.Open 106 } 107 108 // Build payload 109 payload := &models.WebhookPayload{ 110 Ref: ref, 111 Before: oldSha, 112 After: newSha, 113 Repository: repository, 114 Pusher: models.WebhookUser{ 115 Did: pusher, 116 }, 117 } 118 119 return payload, nil 120} 121 122// sendWebhook sends the webhook http request 123func (w *WebhookNotifier) sendWebhook(ctx context.Context, webhook models.Webhook, event string, payload *models.WebhookPayload) { 124 deliveryId := uuid.New().String() 125 126 payloadBytes, err := json.Marshal(payload) 127 if err != nil { 128 w.logger.Error("failed to marshal webhook payload", "webhook_id", webhook.Id, "err", err) 129 return 130 } 131 132 req, err := http.NewRequestWithContext(ctx, "POST", webhook.Url, bytes.NewReader(payloadBytes)) 133 if err != nil { 134 w.logger.Error("failed to create webhook request", "webhook_id", webhook.Id, "err", err) 135 return 136 } 137 138 shortSha := payload.After[:7] 139 140 req.Header.Set("Content-Type", "application/json") 141 req.Header.Set("User-Agent", "Tangled-Hook/"+shortSha) 142 req.Header.Set("X-Tangled-Event", event) 143 req.Header.Set("X-Tangled-Hook-ID", fmt.Sprintf("%d", webhook.Id)) 144 req.Header.Set("X-Tangled-Delivery", deliveryId) 145 req.Header.Set("X-Tangled-Repo", payload.Repository.FullName) 146 147 if webhook.Secret != "" { 148 signature := w.computeSignature(payloadBytes, webhook.Secret) 149 req.Header.Set("X-Tangled-Signature-256", "sha256="+signature) 150 } 151 152 delivery := &models.WebhookDelivery{ 153 WebhookId: webhook.Id, 154 Event: event, 155 DeliveryId: deliveryId, 156 Url: webhook.Url, 157 RequestBody: string(payloadBytes), 158 } 159 160 // retry webhook delivery with exponential backoff 161 retryOpts := []retry.Option{ 162 retry.Attempts(3), 163 retry.Delay(1 * time.Second), 164 retry.MaxDelay(10 * time.Second), 165 retry.DelayType(retry.BackOffDelay), 166 retry.LastErrorOnly(true), 167 retry.OnRetry(func(n uint, err error) { 168 w.logger.Info("retrying webhook delivery", 169 "webhook_id", webhook.Id, 170 "attempt", n+1, 171 "err", err) 172 }), 173 retry.Context(ctx), 174 retry.RetryIf(func(err error) bool { 175 // only retry on network errors or 5xx responses 176 if err != nil { 177 return true 178 } 179 return false 180 }), 181 } 182 183 var resp *http.Response 184 err = retry.Do(func() error { 185 var err error 186 resp, err = w.client.Do(req) 187 if err != nil { 188 return err 189 } 190 191 // retry on 5xx server errors 192 if resp.StatusCode >= 500 { 193 defer resp.Body.Close() 194 return fmt.Errorf("server error: %d", resp.StatusCode) 195 } 196 197 return nil 198 }, retryOpts...) 199 200 if err != nil { 201 w.logger.Error("webhook request failed after retries", "webhook_id", webhook.Id, "err", err) 202 delivery.Success = false 203 delivery.ResponseBody = err.Error() 204 } else { 205 defer resp.Body.Close() 206 207 delivery.ResponseCode = resp.StatusCode 208 delivery.Success = resp.StatusCode >= 200 && resp.StatusCode < 300 209 210 // Read response body (limit to 10KB) 211 bodyBytes, err := io.ReadAll(io.LimitReader(resp.Body, 10*1024)) 212 if err != nil { 213 w.logger.Warn("failed to read webhook response body", "webhook_id", webhook.Id, "err", err) 214 } else { 215 delivery.ResponseBody = string(bodyBytes) 216 } 217 218 if !delivery.Success { 219 w.logger.Warn("webhook delivery failed", 220 "webhook_id", webhook.Id, 221 "status", resp.StatusCode, 222 "url", webhook.Url) 223 } else { 224 w.logger.Info("webhook delivered successfully", 225 "webhook_id", webhook.Id, 226 "url", webhook.Url, 227 "delivery_id", deliveryId) 228 } 229 } 230 231 if err := db.AddWebhookDelivery(w.db, delivery); err != nil { 232 w.logger.Error("failed to record webhook delivery", "webhook_id", webhook.Id, "err", err) 233 } 234} 235 236// computeSignature computes HMAC-SHA256 signature for the payload 237func (w *WebhookNotifier) computeSignature(payload []byte, secret string) string { 238 mac := hmac.New(sha256.New, []byte(secret)) 239 mac.Write(payload) 240 return hex.EncodeToString(mac.Sum(nil)) 241}