Files
2026-08-23 00:48:10 +08:00

171 lines
4.3 KiB
Go

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)
}