Compare commits

2 Commits

Author SHA1 Message Date
226474b434 轮训有问题
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m7s
2026-05-22 16:21:46 +08:00
860f589b9a 错误提示有问题 2026-05-22 15:54:54 +08:00
6 changed files with 47 additions and 22 deletions

View File

@@ -314,10 +314,13 @@ func startPollingInitializer(ctx context.Context, runtime *workerRuntime, appLog
pollingInitializer.StartBackground(ctx) pollingInitializer.StartBackground(ctx)
runtime.pollingConfigMgr.WatchChanges(ctx, func(hadConfigs, hasConfigs bool) { runtime.pollingConfigMgr.WatchChanges(ctx, func(hadConfigs, hasConfigs bool) {
if !hadConfigs && hasConfigs { if hasConfigs {
appLogger.Info("轮询配置从空变为非空,触发队列重新初始化") appLogger.Info("轮询配置已变更,触发队列重新初始化",
zap.Bool("had_configs", hadConfigs))
pollingInitializer.Restart(ctx) pollingInitializer.Restart(ctx)
return
} }
appLogger.Info("轮询配置已清空,跳过队列重新初始化")
}) })
return pollingInitializer return pollingInitializer

View File

@@ -52,6 +52,7 @@ type PollingInitializer struct {
progress initProgress progress initProgress
initCompleted atomic.Bool initCompleted atomic.Bool
restartQueued atomic.Bool
stopChan chan struct{} stopChan chan struct{}
wg sync.WaitGroup wg sync.WaitGroup
@@ -94,11 +95,12 @@ func (p *PollingInitializer) IsCompleted() bool {
return p.initCompleted.Load() return p.initCompleted.Load()
} }
// Restart 重新执行初始化(当轮询配置从空变为非空时调用) // Restart 重新执行初始化,用于配置变更后补齐新增任务类型的分片队列。
// 使用 CAS 确保只有初始化已完成时才能重启,避免并发重入 // 若初始化正在进行中,则记录一次待重启请求,当前轮完成后再自动执行。
func (p *PollingInitializer) Restart(ctx context.Context) { func (p *PollingInitializer) Restart(ctx context.Context) {
if !p.initCompleted.CompareAndSwap(true, false) { if !p.initCompleted.CompareAndSwap(true, false) {
p.logger.Info("轮询初始化仍在进行中,跳过重启") p.restartQueued.Store(true)
p.logger.Info("轮询初始化仍在进行中,已标记完成后重新初始化")
return return
} }
p.setStatus("pending", "") p.setStatus("pending", "")
@@ -200,6 +202,20 @@ func (p *PollingInitializer) run(ctx context.Context) {
p.logger.Info("分片渐进式初始化完成", p.logger.Info("分片渐进式初始化完成",
zap.Int64("total_loaded", snapshot.LoadedCards), zap.Int64("total_loaded", snapshot.LoadedCards),
zap.Duration("duration", time.Since(snapshot.StartTime))) zap.Duration("duration", time.Since(snapshot.StartTime)))
if p.restartQueued.Swap(false) {
select {
case <-ctx.Done():
return
default:
}
if p.initCompleted.CompareAndSwap(true, false) {
p.setStatus("pending", "")
p.wg.Add(1)
go p.run(ctx)
p.logger.Info("检测到初始化期间配置变更,完成后再次初始化")
}
}
} }
// initBatch 使用 Pipeline 将一批卡写入分片队列和缓存 // initBatch 使用 Pipeline 将一批卡写入分片队列和缓存

View File

@@ -155,6 +155,7 @@ func (s *ConfigService) Update(ctx context.Context, id uint, req *dto.UpdatePoll
return nil, errors.Wrap(errors.CodeInternalError, err, "更新轮询配置失败") return nil, errors.Wrap(errors.CodeInternalError, err, "更新轮询配置失败")
} }
s.notifyConfigChanged(ctx, "updated")
return s.toResponse(config), nil return s.toResponse(config), nil
} }

View File

@@ -1,35 +1,37 @@
package alipay package alipay
import ( import (
"fmt"
"github.com/smartwalle/alipay/v3" "github.com/smartwalle/alipay/v3"
"github.com/break/junhong_cmp_fiber/internal/model" "github.com/break/junhong_cmp_fiber/internal/model"
apperrors "github.com/break/junhong_cmp_fiber/pkg/errors"
) )
// NewClientFromConfig 从 WechatConfig 构建支付宝 SDK client。 // NewClientFromConfig 从 WechatConfig 构建支付宝 SDK client。
// 校验 ali_app_id、ali_private_key、ali_public_key 必须非空。 // 校验 ali_app_id、ali_private_key、ali_public_key 必须非空。
// ali_production=true 时使用生产环境,否则使用沙箱环境。 // ali_production=true 时使用生产环境,否则使用沙箱环境。
func NewClientFromConfig(cfg *model.WechatConfig) (*alipay.Client, error) { func NewClientFromConfig(cfg *model.WechatConfig) (*alipay.Client, error) {
if cfg == nil {
return nil, apperrors.New(apperrors.CodeNoPaymentConfig, "支付宝配置缺失:未找到启用的支付配置")
}
if cfg.AliAppID == "" { if cfg.AliAppID == "" {
return nil, fmt.Errorf("支付宝配置缺失ali_app_id 不能为空") return nil, apperrors.New(apperrors.CodeNoPaymentConfig, "支付宝配置缺失ali_app_id 不能为空")
} }
if cfg.AliPrivateKey == "" { if cfg.AliPrivateKey == "" {
return nil, fmt.Errorf("支付宝配置缺失ali_private_key 不能为空") return nil, apperrors.New(apperrors.CodeNoPaymentConfig, "支付宝配置缺失ali_private_key 不能为空")
} }
if cfg.AliPublicKey == "" { if cfg.AliPublicKey == "" {
return nil, fmt.Errorf("支付宝配置缺失ali_public_key 不能为空") return nil, apperrors.New(apperrors.CodeNoPaymentConfig, "支付宝配置缺失ali_public_key 不能为空")
} }
client, err := alipay.New(cfg.AliAppID, cfg.AliPrivateKey, cfg.AliProduction) client, err := alipay.New(cfg.AliAppID, cfg.AliPrivateKey, cfg.AliProduction)
if err != nil { if err != nil {
return nil, fmt.Errorf("创建支付宝客户端失败: %w", err) return nil, apperrors.Wrap(apperrors.CodeNoPaymentConfig, err, "支付宝配置不可用:创建客户端失败")
} }
// 加载支付宝公钥(公钥模式验签) // 加载支付宝公钥(公钥模式验签)
if err := client.LoadAliPayPublicKey(cfg.AliPublicKey); err != nil { if err := client.LoadAliPayPublicKey(cfg.AliPublicKey); err != nil {
return nil, fmt.Errorf("加载支付宝公钥失败: %w", err) return nil, apperrors.Wrap(apperrors.CodeNoPaymentConfig, err, "支付宝配置不可用:加载公钥失败")
} }
return client, nil return client, nil

View File

@@ -2,12 +2,12 @@ package alipay
import ( import (
"context" "context"
"fmt"
"time" "time"
"github.com/smartwalle/alipay/v3" "github.com/smartwalle/alipay/v3"
"github.com/break/junhong_cmp_fiber/internal/model" "github.com/break/junhong_cmp_fiber/internal/model"
apperrors "github.com/break/junhong_cmp_fiber/pkg/errors"
) )
// BuildWapPayURL 生成支付宝手机网站支付 URL。 // BuildWapPayURL 生成支付宝手机网站支付 URL。
@@ -40,7 +40,7 @@ func BuildWapPayURL(ctx context.Context, cfg *model.WechatConfig, payment *model
payURL, err := client.TradeWapPay(param) payURL, err := client.TradeWapPay(param)
if err != nil { if err != nil {
return "", fmt.Errorf("生成支付宝 WAP 支付链接失败: %w", err) return "", apperrors.Wrap(apperrors.CodeNoPaymentConfig, err, "支付宝配置不可用:生成支付链接失败")
} }
return payURL.String(), nil return payURL.String(), nil

View File

@@ -1,6 +1,7 @@
package errors package errors
import ( import (
stderrors "errors"
"runtime/debug" "runtime/debug"
"time" "time"
@@ -51,12 +52,14 @@ func handleError(c *fiber.Ctx, err error, logger *zap.Logger) error {
var message string var message string
var httpStatus int var httpStatus int
switch e := err.(type) { var appErr *AppError
case *AppError: var fiberErr *fiber.Error
switch {
case stderrors.As(err, &appErr):
// 应用自定义错误 // 应用自定义错误
code = e.Code code = appErr.Code
message = e.Message message = appErr.Message
httpStatus = GetHTTPStatus(e.Code) httpStatus = GetHTTPStatus(appErr.Code)
// 记录错误日志(包含完整上下文) // 记录错误日志(包含完整上下文)
logFields := append(errCtx.ToLogFields(), logFields := append(errCtx.ToLogFields(),
@@ -74,16 +77,16 @@ func handleError(c *fiber.Ctx, err error, logger *zap.Logger) error {
safeLogWithLevel(logger, "warn", "客户端错误", logFields...) safeLogWithLevel(logger, "warn", "客户端错误", logFields...)
} }
case *fiber.Error: case stderrors.As(err, &fiberErr):
// Fiber 框架错误 // Fiber 框架错误
httpStatus = e.Code httpStatus = fiberErr.Code
code = mapHTTPStatusToCode(httpStatus) code = mapHTTPStatusToCode(httpStatus)
message = GetMessage(code, "zh") message = GetMessage(code, "zh")
safeLog(logger, "Fiber 框架错误", safeLog(logger, "Fiber 框架错误",
append(errCtx.ToLogFields(), append(errCtx.ToLogFields(),
zap.Int("http_status", httpStatus), zap.Int("http_status", httpStatus),
zap.String("fiber_message", e.Message), zap.String("fiber_message", fiberErr.Message),
)..., )...,
) )