更新前后端代码
This commit is contained in:
@@ -0,0 +1,170 @@
|
||||
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)
|
||||
}
|
||||
Reference in New Issue
Block a user