实现 sse 检查订单,减少请求次数
This commit is contained in:
@@ -1,33 +0,0 @@
|
||||
package events
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"github.com/hibiken/asynq"
|
||||
"log/slog"
|
||||
"time"
|
||||
)
|
||||
|
||||
type RequestLog struct {
|
||||
Type string
|
||||
User int32
|
||||
IP string
|
||||
UA string
|
||||
Method string
|
||||
Path string
|
||||
Status int
|
||||
Error string
|
||||
Latency time.Duration
|
||||
Time time.Time
|
||||
}
|
||||
|
||||
func NewRequestLog(data RequestLog) *asynq.Task {
|
||||
var rs, err = json.Marshal(data)
|
||||
if err != nil {
|
||||
slog.Error("日志数据序列化失败", slog.Any("err", err), slog.Any("data", data))
|
||||
return nil
|
||||
}
|
||||
return asynq.NewTask("logs:request", rs)
|
||||
}
|
||||
|
||||
type LoginLog struct {
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package globals
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
|
||||
"github.com/go-playground/locales/zh"
|
||||
@@ -19,22 +20,32 @@ type ValidatorClient struct {
|
||||
translator ut.Translator
|
||||
}
|
||||
|
||||
func (v *ValidatorClient) Validate(c *fiber.Ctx, data any) error {
|
||||
|
||||
func (v *ValidatorClient) ParseBody(c *fiber.Ctx, data any) error {
|
||||
if err := c.BodyParser(data); err != nil {
|
||||
return err
|
||||
}
|
||||
return validate(data)
|
||||
}
|
||||
|
||||
if errs := v.validator.Struct(data); errs != nil {
|
||||
func (v *ValidatorClient) ParseQuery(c *fiber.Ctx, data any) error {
|
||||
if err := c.QueryParser(data); err != nil {
|
||||
return err
|
||||
}
|
||||
return validate(data)
|
||||
}
|
||||
|
||||
func validate(data any) error {
|
||||
if errs := Validator.validator.Struct(data); errs != nil {
|
||||
var sb = strings.Builder{}
|
||||
var typeErrs validator.ValidationErrors
|
||||
errors.As(errs, &typeErrs)
|
||||
for i, err := range typeErrs {
|
||||
sb.WriteString(err.Translate(v.translator))
|
||||
sb.WriteString(err.Translate(Validator.translator))
|
||||
if i < len(typeErrs)-1 {
|
||||
sb.WriteString("\n")
|
||||
}
|
||||
}
|
||||
slog.Debug("请求参数验证失败", "msg", sb.String())
|
||||
return fiber.NewError(fiber.StatusBadRequest, sb.String())
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -113,7 +113,7 @@ func CreateChannel(c *fiber.Ctx) error {
|
||||
|
||||
// 解析参数
|
||||
req := new(CreateChannelReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return core.NewBizErr("解析参数失败", err)
|
||||
}
|
||||
|
||||
|
||||
@@ -147,7 +147,7 @@ func ListResourceLong(c *fiber.Ctx) error {
|
||||
do.Where(q.ResourceLong.As(q.Resource.Long.Name()).Expire.Lte(*req.ExpireBefore))
|
||||
}
|
||||
|
||||
resource, err := q.Resource.Debug().Where(do).
|
||||
resource, err := q.Resource.Where(do).
|
||||
Joins(q.Resource.Long).
|
||||
Order(q.Resource.CreatedAt.Desc()).
|
||||
Offset(req.GetOffset()).
|
||||
@@ -354,7 +354,7 @@ func StatisticResourceUsage(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
var req = new(StatisticResourceUsageReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -416,7 +416,7 @@ func CreateResource(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
var req = new(CreateResourceReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -438,7 +438,7 @@ func ResourcePrice(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
var req = new(CreateResourceReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
@@ -1,15 +1,20 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"platform/pkg/env"
|
||||
"platform/web/auth"
|
||||
"platform/web/core"
|
||||
g "platform/web/globals"
|
||||
m "platform/web/models"
|
||||
s "platform/web/services"
|
||||
"reflect"
|
||||
"time"
|
||||
|
||||
"github.com/gofiber/fiber/v2"
|
||||
"github.com/valyala/fasthttp"
|
||||
)
|
||||
|
||||
type TradeCreateReq struct {
|
||||
@@ -33,7 +38,7 @@ func TradeCreate(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
req := new(TradeCreateReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -76,7 +81,7 @@ func TradeComplete(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
req := new(TradeCompleteReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -102,7 +107,7 @@ func TradeCancel(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
req := new(TradeCancelReq)
|
||||
if err := g.Validator.Validate(c, req); err != nil {
|
||||
if err := g.Validator.ParseBody(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -115,3 +120,54 @@ func TradeCancel(c *fiber.Ctx) error {
|
||||
|
||||
return c.SendStatus(fiber.StatusNoContent)
|
||||
}
|
||||
|
||||
type TradeCheckReq struct {
|
||||
s.ModifyTradeData
|
||||
}
|
||||
|
||||
func TradeCheck(c *fiber.Ctx) error {
|
||||
// 解析请求参数
|
||||
req := new(TradeCheckReq)
|
||||
if err := g.Validator.ParseQuery(c, req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.Set(fiber.HeaderContentType, "text/event-stream")
|
||||
c.Set(fiber.HeaderCacheControl, "no-cache")
|
||||
c.Set(fiber.HeaderConnection, "keep-alive")
|
||||
c.Set(fiber.HeaderTransferEncoding, "chunked")
|
||||
c.Context().SetBodyStreamWriter(fasthttp.StreamWriter(func(w *bufio.Writer) {
|
||||
|
||||
expire := env.TradeExpire
|
||||
interval := 5
|
||||
for range expire / interval {
|
||||
// 检查订单状态
|
||||
result, err := s.Trade.CheckTrade(&req.ModifyTradeData)
|
||||
if err != nil {
|
||||
slog.Error("检查订单状态失败", "trade_no", req.TradeNo, "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
// 写入订单状态
|
||||
_, err = fmt.Fprintf(w, "data: %d\n\n", result.Status)
|
||||
if err != nil {
|
||||
slog.Error("写入订单状态失败", "trade_no", req.TradeNo, "error", err)
|
||||
return
|
||||
}
|
||||
err = w.Flush()
|
||||
if err != nil {
|
||||
slog.Error("刷新缓冲区失败", "trade_no", req.TradeNo, "error", err, "errType", reflect.TypeOf(err))
|
||||
return
|
||||
}
|
||||
|
||||
// 当订单离开支付状态后结束查询
|
||||
if result.Status != m.TradeStatusPending {
|
||||
return
|
||||
}
|
||||
|
||||
time.Sleep(time.Duration(interval) * time.Second)
|
||||
}
|
||||
}))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -85,7 +85,7 @@ func CreateWhitelist(c *fiber.Ctx) error {
|
||||
|
||||
// 解析请求参数
|
||||
req := new(CreateWhitelistReq)
|
||||
err = g.Validator.Validate(c, req)
|
||||
err = g.Validator.ParseBody(c, req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
|
||||
"github.com/gofiber/contrib/otelfiber/v2"
|
||||
"github.com/gofiber/fiber/v2"
|
||||
"github.com/gofiber/fiber/v2/middleware/cors"
|
||||
"github.com/gofiber/fiber/v2/middleware/logger"
|
||||
"github.com/gofiber/fiber/v2/middleware/recover"
|
||||
"github.com/gofiber/fiber/v2/middleware/requestid"
|
||||
@@ -19,9 +20,6 @@ func ApplyMiddlewares(app *fiber.App) {
|
||||
EnableStackTrace: true,
|
||||
}))
|
||||
|
||||
// metric
|
||||
app.Use(otelfiber.Middleware())
|
||||
|
||||
// logger
|
||||
app.Use(logger.New(logger.Config{
|
||||
Next: func(c *fiber.Ctx) bool {
|
||||
@@ -29,6 +27,9 @@ func ApplyMiddlewares(app *fiber.App) {
|
||||
},
|
||||
}))
|
||||
|
||||
// metric
|
||||
app.Use(otelfiber.Middleware())
|
||||
|
||||
// request id
|
||||
app.Use(requestid.New(requestid.Config{
|
||||
Generator: func() string {
|
||||
@@ -37,6 +38,9 @@ func ApplyMiddlewares(app *fiber.App) {
|
||||
},
|
||||
}))
|
||||
|
||||
// cors
|
||||
app.Use(cors.New())
|
||||
|
||||
// authenticate
|
||||
app.Use(auth.Authenticate())
|
||||
}
|
||||
|
||||
@@ -52,6 +52,7 @@ func ApplyRouters(app *fiber.App) {
|
||||
trade.Post("/create", handlers.TradeCreate)
|
||||
trade.Post("/complete", handlers.TradeComplete)
|
||||
trade.Post("/cancel", handlers.TradeCancel)
|
||||
trade.Get("/check", handlers.TradeCheck)
|
||||
|
||||
// 账单
|
||||
bill := api.Group("/bill")
|
||||
|
||||
@@ -283,7 +283,7 @@ func (s *channelBaiyinService) CreateChannels(source netip.Addr, resourceId int3
|
||||
}
|
||||
} else {
|
||||
bytes, _ := json.Marshal(configs)
|
||||
slog.Debug("提交代理端口配置", "config", string(bytes))
|
||||
slog.Debug("提交代理端口配置", "proxy", proxy.IP.String(), "config", string(bytes))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -348,7 +348,7 @@ func (s *channelBaiyinService) RemoveChannels(batch string, ids []int32) error {
|
||||
}
|
||||
} else {
|
||||
bytes, _ := json.Marshal(configs)
|
||||
slog.Debug("清除代理端口配置", "config", string(bytes))
|
||||
slog.Debug("清除代理端口配置", "proxy", ip, "config", string(bytes))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -668,7 +668,7 @@ type CreateTradeResult struct {
|
||||
}
|
||||
|
||||
type ModifyTradeData struct {
|
||||
TradeNo string `json:"trade_no" validate:"required"`
|
||||
TradeNo string `json:"trade_no" query:"trade_no" validate:"required"`
|
||||
Method m.TradeMethod `json:"method" validate:"required"`
|
||||
}
|
||||
|
||||
|
||||
@@ -141,7 +141,7 @@ func HandleFlushGateway(_ context.Context, task *asynq.Task) error {
|
||||
}
|
||||
} else {
|
||||
bytes, _ := json.Marshal(configs)
|
||||
slog.Debug("更新代理后备配置", "config", string(bytes))
|
||||
slog.Debug("更新代理后备配置", "proxy", proxy.IP.String(), "config", string(bytes))
|
||||
}
|
||||
|
||||
_, err := q.Proxy.
|
||||
|
||||
Reference in New Issue
Block a user