2026-09-12 15:32:27 +08:00
|
|
|
// Package gateway 提供消息网关的 HTTP 接口。
|
|
|
|
|
//
|
|
|
|
|
// 数据库读写都发生在这一层,真正的投递逻辑在 internal/gateway 里,
|
|
|
|
|
// 这里只负责鉴权、取数、组装投递任务。
|
2026-08-30 21:09:26 +08:00
|
|
|
package gateway
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"net/http"
|
|
|
|
|
"strings"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"github.com/gin-gonic/gin"
|
|
|
|
|
"github.com/ssdomei232/goodBaby/api/response"
|
|
|
|
|
"github.com/ssdomei232/goodBaby/api/user"
|
|
|
|
|
"github.com/ssdomei232/goodBaby/handler/db"
|
2026-09-12 15:32:27 +08:00
|
|
|
gatewaycore "github.com/ssdomei232/goodBaby/internal/gateway"
|
2026-08-30 21:09:26 +08:00
|
|
|
"github.com/ssdomei232/goodBaby/internal/retry"
|
|
|
|
|
"github.com/ssdomei232/goodBaby/model"
|
2026-09-12 15:32:27 +08:00
|
|
|
"gorm.io/gorm"
|
2026-08-30 21:09:26 +08:00
|
|
|
)
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
// webhookRequest 外部系统投递消息的请求体
|
2026-08-30 21:09:26 +08:00
|
|
|
type webhookRequest struct {
|
|
|
|
|
Message string `json:"message"`
|
|
|
|
|
Title string `json:"title"`
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
// HandleList 获取当前用户的所有消息网关
|
2026-08-30 21:09:26 +08:00
|
|
|
func HandleList(c *gin.Context) {
|
2026-09-12 15:32:27 +08:00
|
|
|
userInfo, err := user.GetUserInfoByGinCtx(c)
|
2026-08-30 21:09:26 +08:00
|
|
|
if err != nil {
|
|
|
|
|
response.Unauthorized(c, "未登录")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
dbConn, err := db.GetGormDB()
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, "获取网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
items := []model.MessageGateway{}
|
|
|
|
|
if err := dbConn.Where("uid = ?", userInfo.ID).Order("id DESC").Find(&items).Error; err != nil {
|
2026-08-30 21:09:26 +08:00
|
|
|
response.ServerError(c, "获取网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
response.OK(c, items)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
// HandleCreate 创建消息网关
|
2026-08-30 21:09:26 +08:00
|
|
|
func HandleCreate(c *gin.Context) {
|
2026-09-12 15:32:27 +08:00
|
|
|
userInfo, err := user.GetUserInfoByGinCtx(c)
|
2026-08-30 21:09:26 +08:00
|
|
|
if err != nil {
|
|
|
|
|
response.Unauthorized(c, "未登录")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
var req model.MessageGatewayRequest
|
2026-09-12 15:32:27 +08:00
|
|
|
if err := c.ShouldBindJSON(&req); err != nil {
|
|
|
|
|
response.BadRequest(c, "输入参数错误")
|
2026-08-30 21:09:26 +08:00
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
req.Name = strings.TrimSpace(req.Name)
|
|
|
|
|
if err := req.Validate(); err != nil {
|
|
|
|
|
response.FromError(c, err, "创建网关失败")
|
2026-08-30 21:09:26 +08:00
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
// 不填类型时用默认网关,填了就必须是已注册的类型
|
|
|
|
|
if req.Type == "" {
|
|
|
|
|
req.Type = model.GatewayTypeWebhook
|
|
|
|
|
}
|
|
|
|
|
if _, ok := gatewaycore.InitGatewayRegistry().Resolve(req.Type); !ok {
|
|
|
|
|
response.BadRequest(c, "不支持的消息网关类型: "+req.Type)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
token, err := gatewaycore.NewToken()
|
2026-08-30 21:09:26 +08:00
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, "生成网关 Token 失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
dbConn, err := db.GetGormDB()
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, "创建网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
item := model.MessageGateway{
|
|
|
|
|
UID: userInfo.ID,
|
|
|
|
|
Name: req.Name,
|
|
|
|
|
Type: req.Type,
|
|
|
|
|
Token: token,
|
|
|
|
|
CreateAt: time.Now().Unix(),
|
|
|
|
|
}
|
2026-08-30 21:09:26 +08:00
|
|
|
if err := dbConn.Create(&item).Error; err != nil {
|
|
|
|
|
response.ServerError(c, "创建网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
response.OK(c, item)
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
// HandleDelete 删除消息网关,绑定在它下面的规则一并删除
|
2026-08-30 21:09:26 +08:00
|
|
|
func HandleDelete(c *gin.Context) {
|
2026-09-12 15:32:27 +08:00
|
|
|
userInfo, err := user.GetUserInfoByGinCtx(c)
|
2026-08-30 21:09:26 +08:00
|
|
|
if err != nil {
|
|
|
|
|
response.Unauthorized(c, "未登录")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
gatewayID, err := parseID(c.Param("gatewayID"))
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.BadRequest(c, "网关 ID 格式错误")
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
dbConn, err := db.GetGormDB()
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, "删除网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
err = dbConn.Transaction(func(tx *gorm.DB) error {
|
|
|
|
|
if err := tx.Where("gateway_id = ? AND uid = ?", gatewayID, userInfo.ID).
|
|
|
|
|
Delete(&model.GatewayRule{}).Error; err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
return tx.Where("id = ? AND uid = ?", gatewayID, userInfo.ID).
|
|
|
|
|
Delete(&model.MessageGateway{}).Error
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
2026-08-30 21:09:26 +08:00
|
|
|
response.ServerError(c, "删除网关失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
response.OK(c, "网关已删除")
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
// HandleWebhook 接收外部系统投递的消息,触发绑定在该网关上的规则
|
2026-08-30 21:09:26 +08:00
|
|
|
func HandleWebhook(c *gin.Context) {
|
2026-09-12 15:32:27 +08:00
|
|
|
var req webhookRequest
|
|
|
|
|
if err := c.ShouldBindJSON(&req); err != nil || strings.TrimSpace(req.Message) == "" {
|
|
|
|
|
response.BadRequest(c, "message 不能为空")
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-30 21:09:26 +08:00
|
|
|
dbConn, err := db.GetGormDB()
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, "网关不可用")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
var target model.MessageGateway
|
|
|
|
|
if err := dbConn.Where("token = ?", c.Param("token")).First(&target).Error; err != nil {
|
2026-08-30 21:09:26 +08:00
|
|
|
response.Fail(c, http.StatusNotFound, "网关不存在")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
rules := []model.GatewayRule{}
|
|
|
|
|
if err := dbConn.Where("uid = ? AND gateway_id = ? AND enabled = ?", target.UID, target.ID, true).
|
|
|
|
|
Find(&rules).Error; err != nil {
|
2026-08-30 21:09:26 +08:00
|
|
|
response.ServerError(c, "读取网关规则失败")
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
// 投递需要在 HTTP 请求内返回结果,因此用较短的超时
|
2026-08-30 21:09:26 +08:00
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), retry.TestTimeout)
|
|
|
|
|
defer cancel()
|
|
|
|
|
|
2026-09-12 15:32:27 +08:00
|
|
|
result, err := gatewaycore.InitGatewayRegistry().Deliver(ctx, target.Type, &gatewaycore.Task{
|
|
|
|
|
Gateway: &target,
|
|
|
|
|
Rules: rules,
|
|
|
|
|
Message: gatewaycore.Message{
|
|
|
|
|
Title: strings.TrimSpace(req.Title),
|
|
|
|
|
Content: strings.TrimSpace(req.Message),
|
|
|
|
|
},
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
response.ServerError(c, err.Error())
|
|
|
|
|
return
|
2026-08-30 21:09:26 +08:00
|
|
|
}
|
2026-09-12 15:32:27 +08:00
|
|
|
|
|
|
|
|
response.OK(c, result)
|
2026-08-30 21:09:26 +08:00
|
|
|
}
|