回调优化为单地址 upsert 配置并支持重置密钥
- 回调订阅改为单记录 upsert:保存即覆盖(URL/事件/状态),自动停用其他订阅,无新增/删除 - 新增重置密钥:保存时传 rotate_secret=true 重新生成密钥并返回一次,解决密钥丢失无法找回 - 删除不再使用的 PATCH /callbacks/:id/status 路由及对应 handler/service 死代码 - 前端回调页改为内联表单(URL/事件/状态),新增重置密钥按钮(确认后展示新密钥一次) - 补充单地址复用、disabled 不投递、密钥轮换测试
This commit is contained in:
@@ -21,27 +21,28 @@ import (
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
// 回调推送策略(参考微信/支付宝通知机制):
|
||||
// 首次推送失败后按固定序列退避重试,默认共 16 次尝试,之后标记 failed 不再推送。
|
||||
// 序列与总次数均可通过环境变量覆盖(CALLBACK_RETRY_SCHEDULE / CALLBACK_MAX_ATTEMPTS)。
|
||||
var defaultCallbackRetrySchedule = []time.Duration{
|
||||
15 * time.Second, // 第 2 次
|
||||
15 * time.Second, // 第 3 次
|
||||
30 * time.Second, // 第 4 次
|
||||
3 * time.Minute, // 第 5 次
|
||||
10 * time.Minute, // 第 6 次
|
||||
20 * time.Minute, // 第 7 次
|
||||
30 * time.Minute, // 第 8 次
|
||||
30 * time.Minute, // 第 9 次
|
||||
30 * time.Minute, // 第 10 次
|
||||
60 * time.Minute, // 第 11 次
|
||||
3 * time.Hour, // 第 12 次
|
||||
3 * time.Hour, // 第 13 次
|
||||
3 * time.Hour, // 第 14 次
|
||||
6 * time.Hour, // 第 15 次
|
||||
6 * time.Hour, // 第 16 次
|
||||
15 * time.Second, // 第 2 次
|
||||
15 * time.Second, // 第 3 次
|
||||
30 * time.Second, // 第 4 次
|
||||
3 * time.Minute, // 第 5 次
|
||||
10 * time.Minute, // 第 6 次
|
||||
20 * time.Minute, // 第 7 次
|
||||
30 * time.Minute, // 第 8 次
|
||||
30 * time.Minute, // 第 9 次
|
||||
30 * time.Minute, // 第 10 次
|
||||
60 * time.Minute, // 第 11 次
|
||||
3 * time.Hour, // 第 12 次
|
||||
3 * time.Hour, // 第 13 次
|
||||
3 * time.Hour, // 第 14 次
|
||||
6 * time.Hour, // 第 15 次
|
||||
6 * time.Hour, // 第 16 次
|
||||
}
|
||||
|
||||
// CallbackConfig 回调推送策略配置。
|
||||
@@ -54,10 +55,10 @@ type CallbackConfig struct {
|
||||
|
||||
// CallbackService 以数据库 outbox 方式管理回调,进程重启不会丢失待发送事件。
|
||||
type CallbackService struct {
|
||||
db *gorm.DB
|
||||
codec *SecretCodec
|
||||
httpClient *http.Client
|
||||
maxAttempts int
|
||||
db *gorm.DB
|
||||
codec *SecretCodec
|
||||
httpClient *http.Client
|
||||
maxAttempts int
|
||||
retrySchedule []time.Duration
|
||||
}
|
||||
|
||||
@@ -103,6 +104,9 @@ type CreateCallbackInput struct {
|
||||
Name string
|
||||
URL string
|
||||
Events string
|
||||
Status string
|
||||
// RotateSecret 重置回调密钥:重新生成 secret(旧密钥立即失效),返回新 secret 仅此一次。
|
||||
RotateSecret bool
|
||||
}
|
||||
|
||||
type CallbackCredential struct {
|
||||
@@ -114,7 +118,7 @@ func (s *CallbackService) CreateSubscription(merchantID uint, in CreateCallbackI
|
||||
in.Name = strings.TrimSpace(in.Name)
|
||||
in.URL = strings.TrimSpace(in.URL)
|
||||
if in.Name == "" {
|
||||
return nil, errors.New("回调名称不能为空")
|
||||
in.Name = "默认回调"
|
||||
}
|
||||
if err := validateCallbackURL(in.URL); err != nil {
|
||||
return nil, err
|
||||
@@ -123,41 +127,69 @@ func (s *CallbackService) CreateSubscription(merchantID uint, in CreateCallbackI
|
||||
if len(events) == 0 {
|
||||
return nil, errors.New("至少订阅一个事件")
|
||||
}
|
||||
// 每个商户仅允许一个回调地址:存在 active 订阅时拒绝创建,先停用旧的再新建。
|
||||
var activeCount int64
|
||||
if err := s.db.Model(&model.CallbackSubscription{}).
|
||||
Where("merchant_id = ? AND status = ?", merchantID, model.CallbackStatusActive).
|
||||
Count(&activeCount).Error; err != nil {
|
||||
return nil, err
|
||||
if in.Status == "" {
|
||||
in.Status = model.CallbackStatusActive
|
||||
}
|
||||
if activeCount > 0 {
|
||||
return nil, errors.New("每个商户仅支持一个回调地址,请先停用当前订阅再创建")
|
||||
if in.Status != model.CallbackStatusActive && in.Status != model.CallbackStatusDisabled {
|
||||
return nil, errors.New("无效的回调状态")
|
||||
}
|
||||
secret, err := randomToken("cb_", 32)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ciphertext, err := s.codec.Encrypt(secret)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
subscription := &model.CallbackSubscription{
|
||||
MerchantID: merchantID,
|
||||
Name: in.Name,
|
||||
URL: in.URL,
|
||||
Events: strings.Join(events, ","),
|
||||
SecretCiphertext: ciphertext,
|
||||
Status: model.CallbackStatusActive,
|
||||
}
|
||||
err = s.db.Transaction(func(tx *gorm.DB) error {
|
||||
|
||||
subscription := &model.CallbackSubscription{}
|
||||
secret := ""
|
||||
err := s.db.Transaction(func(tx *gorm.DB) error {
|
||||
var merchant model.Merchant
|
||||
if err := tx.Where("id = ? AND status = ?", merchantID, model.MerchantStatusActive).First(&merchant).Error; err != nil {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("id = ? AND status = ?", merchantID, model.MerchantStatusActive).
|
||||
First(&merchant).Error; err != nil {
|
||||
return errors.New("商户不存在或已禁用")
|
||||
}
|
||||
if err := tx.Create(subscription).Error; err != nil {
|
||||
|
||||
err := tx.Where("merchant_id = ?", merchantID).Order("id DESC").First(subscription).Error
|
||||
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return err
|
||||
}
|
||||
return writeAudit(tx, &merchantID, &actorUserID, nil, "callback_subscription.create", "callback_subscription", fmt.Sprint(subscription.ID), map[string]string{"url": subscription.URL})
|
||||
action := "callback_subscription.update"
|
||||
rotate := in.RotateSecret
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
rotate = true // 新建时必然生成新密钥
|
||||
action = "callback_subscription.create"
|
||||
}
|
||||
if rotate {
|
||||
var tokenErr error
|
||||
secret, tokenErr = randomToken("cb_", 32)
|
||||
if tokenErr != nil {
|
||||
return tokenErr
|
||||
}
|
||||
ciphertext, encryptErr := s.codec.Encrypt(secret)
|
||||
if encryptErr != nil {
|
||||
return encryptErr
|
||||
}
|
||||
if subscription.ID == 0 {
|
||||
*subscription = model.CallbackSubscription{
|
||||
MerchantID: merchantID,
|
||||
SecretCiphertext: ciphertext,
|
||||
}
|
||||
} else {
|
||||
subscription.SecretCiphertext = ciphertext
|
||||
}
|
||||
}
|
||||
subscription.Name = in.Name
|
||||
subscription.URL = in.URL
|
||||
subscription.Events = strings.Join(events, ",")
|
||||
subscription.Status = in.Status
|
||||
if subscription.ID == 0 {
|
||||
if err := tx.Create(subscription).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
} else if err := tx.Save(subscription).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Model(&model.CallbackSubscription{}).
|
||||
Where("merchant_id = ? AND id <> ? AND status = ?", merchantID, subscription.ID, model.CallbackStatusActive).
|
||||
Update("status", model.CallbackStatusDisabled).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return writeAudit(tx, &merchantID, &actorUserID, nil, action, "callback_subscription", fmt.Sprint(subscription.ID), map[string]string{"url": subscription.URL, "status": subscription.Status})
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -165,65 +197,57 @@ func (s *CallbackService) CreateSubscription(merchantID uint, in CreateCallbackI
|
||||
return &CallbackCredential{Subscription: subscription, Secret: secret}, nil
|
||||
}
|
||||
|
||||
func (s *CallbackService) ListSubscriptions(merchantID uint) ([]model.CallbackSubscription, error) {
|
||||
var subscriptions []model.CallbackSubscription
|
||||
err := s.db.Where("merchant_id = ?", merchantID).Order("id DESC").Find(&subscriptions).Error
|
||||
return subscriptions, err
|
||||
}
|
||||
|
||||
func (s *CallbackService) UpdateSubscriptionStatus(merchantID, id uint, status string, actorUserID uint) error {
|
||||
if status != model.CallbackStatusActive && status != model.CallbackStatusDisabled {
|
||||
return errors.New("无效的回调状态")
|
||||
func (s *CallbackService) GetSubscription(merchantID uint) (*model.CallbackSubscription, error) {
|
||||
var subscription model.CallbackSubscription
|
||||
err := s.db.Where("merchant_id = ?", merchantID).Order("id DESC").First(&subscription).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, nil
|
||||
}
|
||||
return s.db.Transaction(func(tx *gorm.DB) error {
|
||||
result := tx.Model(&model.CallbackSubscription{}).
|
||||
Where("id = ? AND merchant_id = ?", id, merchantID).
|
||||
Update("status", status)
|
||||
if result.Error != nil {
|
||||
return result.Error
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return errors.New("回调订阅不存在")
|
||||
}
|
||||
return writeAudit(tx, &merchantID, &actorUserID, nil, "callback_subscription.status.update", "callback_subscription", fmt.Sprint(id), map[string]string{"status": status})
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &subscription, nil
|
||||
}
|
||||
|
||||
// Enqueue 在调用方事务中写入回调 outbox,只有订单事务成功才会发送事件。
|
||||
func (s *CallbackService) Enqueue(tx *gorm.DB, merchantID uint, event string, data interface{}) error {
|
||||
var subscriptions []model.CallbackSubscription
|
||||
if err := tx.Where("merchant_id = ? AND status = ?", merchantID, model.CallbackStatusActive).Find(&subscriptions).Error; err != nil {
|
||||
var subscription model.CallbackSubscription
|
||||
err := tx.Where("merchant_id = ?", merchantID).
|
||||
Order("id DESC").
|
||||
First(&subscription).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
now := time.Now()
|
||||
for _, subscription := range subscriptions {
|
||||
if !subscribesTo(subscription.Events, event) {
|
||||
continue
|
||||
}
|
||||
eventID := uuid.NewString()
|
||||
payload, err := json.Marshal(map[string]interface{}{
|
||||
"event_id": eventID,
|
||||
"event": event,
|
||||
"occurred_at": timeutil.FormatAPITime(now),
|
||||
"data": data,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
delivery := model.CallbackDelivery{
|
||||
MerchantID: merchantID,
|
||||
CallbackSubscriptionID: subscription.ID,
|
||||
EventID: eventID,
|
||||
Event: event,
|
||||
Payload: string(payload),
|
||||
Status: model.CallbackDeliveryPending,
|
||||
NextAttemptAt: now,
|
||||
}
|
||||
if err := tx.Create(&delivery).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if subscription.Status != model.CallbackStatusActive {
|
||||
return nil
|
||||
}
|
||||
return nil
|
||||
if !subscribesTo(subscription.Events, event) {
|
||||
return nil
|
||||
}
|
||||
now := time.Now()
|
||||
eventID := uuid.NewString()
|
||||
payload, err := json.Marshal(map[string]interface{}{
|
||||
"event_id": eventID,
|
||||
"event": event,
|
||||
"occurred_at": timeutil.FormatAPITime(now),
|
||||
"data": data,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
delivery := model.CallbackDelivery{
|
||||
MerchantID: merchantID,
|
||||
CallbackSubscriptionID: subscription.ID,
|
||||
EventID: eventID,
|
||||
Event: event,
|
||||
Payload: string(payload),
|
||||
Status: model.CallbackDeliveryPending,
|
||||
NextAttemptAt: now,
|
||||
}
|
||||
return tx.Create(&delivery).Error
|
||||
}
|
||||
|
||||
// DispatchDue 执行一批可发送的 outbox 记录。返回成功/失败尝试数,便于日志监控。
|
||||
|
||||
Reference in New Issue
Block a user