package service import ( "bytes" "encoding/json" "fmt" "net/http" "strings" "time" "filestoragesystem/internal/config" "filestoragesystem/internal/model" "filestoragesystem/internal/repository" "filestoragesystem/internal/utils" "filestoragesystem/pkg/apperr" ) // WebhookService Webhook服务 type WebhookService struct { cfg *config.Config webhookRepo *repository.WebhookRepo settingRepo *repository.SettingRepo } // Create 创建webhook func (s *WebhookService) Create(userID uint, projectID *uint, url, secret, events string, status int8) (*model.Webhook, error) { if !strings.HasPrefix(url, "http://") && !strings.HasPrefix(url, "https://") { return nil, fmt.Errorf("回调地址必须以http://或https://开头") } events = normalizeEvents(events) if events == "" { return nil, fmt.Errorf("至少订阅一个事件") } if projectID != nil && *projectID > 0 { // 校验项目归属 var count int64 s.webhookRepo.DB.Model(&model.Project{}).Where("id = ? AND user_id = ?", *projectID, userID).Count(&count) if count == 0 { return nil, apperr.ErrForbidden } } if secret == "" { secret = utils.RandomKey(32) } if status == 0 { status = 1 } w := &model.Webhook{ UserID: userID, ProjectID: projectID, URL: url, Secret: secret, Events: events, Status: status, } if err := s.webhookRepo.Create(w); err != nil { return nil, err } return w, nil } // List 用户webhook列表 func (s *WebhookService) List(userID uint) ([]model.Webhook, error) { return s.webhookRepo.ListByUser(userID) } // Update 更新webhook func (s *WebhookService) Update(userID, id uint, url, events string, status int8) (*model.Webhook, error) { w, err := s.webhookRepo.FindByID(id) if err != nil || w.UserID != userID { return nil, apperr.ErrNotFound } if url != "" { if !strings.HasPrefix(url, "http://") && !strings.HasPrefix(url, "https://") { return nil, fmt.Errorf("回调地址必须以http://或https://开头") } w.URL = url } if events != "" { events = normalizeEvents(events) if events == "" { return nil, fmt.Errorf("至少订阅一个事件") } w.Events = events } if status == 0 || status == 1 { w.Status = status } if err := s.webhookRepo.Update(w); err != nil { return nil, err } return w, nil } // Delete 删除webhook func (s *WebhookService) Delete(userID, id uint) error { if err := s.webhookRepo.Delete(id, userID); err != nil { return apperr.ErrNotFound } return nil } func normalizeEvents(events string) string { seen := make(map[string]bool) parts := make([]string, 0, 4) for _, e := range strings.Split(events, ",") { e = strings.TrimSpace(e) if e == "" || seen[e] { continue } seen[e] = true parts = append(parts, e) } return strings.Join(parts, ",") } func matchEvent(events, event string) bool { for _, e := range strings.Split(events, ",") { if strings.TrimSpace(e) == event { return true } } return false } // Dispatch 异步派发事件回调 func (s *WebhookService) Dispatch(userID, projectID uint, event string, payload map[string]interface{}) { if v, err := s.settingRepo.GetValue("webhook_enabled"); err == nil && v != "true" { return } hooks, err := s.webhookRepo.ListActive(userID, projectID) if err != nil { return } if len(hooks) == 0 { return } body, _ := json.Marshal(map[string]interface{}{ "event": event, "timestamp": time.Now().Unix(), "data": payload, }) for _, hook := range hooks { if !matchEvent(hook.Events, event) { continue } go s.deliver(hook, body) } } func (s *WebhookService) deliver(hook model.Webhook, body []byte) { client := &http.Client{Timeout: 10 * time.Second} req, err := http.NewRequest(http.MethodPost, hook.URL, bytes.NewReader(body)) if err != nil { utils.Warn("webhook", "构建回调请求失败 url=%s err=%v", hook.URL, err) return } req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Webhook-Event", "filestoragesystem") req.Header.Set("X-Webhook-Signature", utils.HMACSHA256(hook.Secret, string(body))) resp, err := client.Do(req) if err != nil { utils.Warn("webhook", "回调发送失败 url=%s err=%v", hook.URL, err) } else { _ = resp.Body.Close() if resp.StatusCode >= 300 { utils.Warn("webhook", "回调响应异常 url=%s status=%d", hook.URL, resp.StatusCode) } } _ = s.webhookRepo.Touch(hook.ID) }