This repository has no description
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}