fedad36089
P01:在 RefundOrder 事务内、lockCommissionSource 之前,先扫描 system_logs 是否已存在该 order_no 的 333 (CommissionTypeRefund) 日志。命中即返回 OrderAlreadyRefunded,不写日志、不动 commission、不动 order.status, 堵住「同一订单被运营人重复退款 → 邀请人佣金被多次扣减」的写入路径。 P02:guard「已退款订单(status=6 + 333 日志)被重新激活」入口。 OrderStatusClaimed(6) 与 orderStatusRefunded(6) 共用同一枚举值, stuckOrderRecovery 把 10 分钟前的 status=6 当成「卡住的 claim」重置回 5 并重新入队 activate,进而让管理员可二次触发退款。新增 logmodel HasRefundCommissionLog helper: - queue/logic/order/stuckOrderRecoveryLogic: 跳过已有 333 日志的订单。 - queue/logic/order/activateOrderLogic.releaseClaim: 同一守卫(防御性)。 新增 internal/model/log/refund.go + 单元测试覆盖 6 个分支(命中、未命中、 子串误判、非法 JSON、空 order_no、DB 错误)。 新增 refundOrderLogic_test.go 覆盖:333 日志已存在 → 直接 rollback、 status==6 短路、status 非 2/5 短路;用 sqlmock 严格断言不再触发 commission 锁/更新/插入。 不做范围: - 不动 activateOrderLogic.calculateCommission(D03.3 单独立项)。 - 不动 OrderStatusClaimed(6) 与 OrderStatusRefunded 枚举值。 - 不动用户 34456 余额数据(D02 待架构师另行决策)。 Co-authored-by: multica-agent <github@multica.ai>
125 lines
4.2 KiB
Go
125 lines
4.2 KiB
Go
package orderLogic
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"time"
|
|
|
|
"github.com/hibiken/asynq"
|
|
logmodel "github.com/perfect-panel/server/internal/model/log"
|
|
"github.com/perfect-panel/server/internal/model/order"
|
|
"github.com/perfect-panel/server/internal/svc"
|
|
"github.com/perfect-panel/server/pkg/logger"
|
|
queueTypes "github.com/perfect-panel/server/queue/types"
|
|
)
|
|
|
|
// StuckOrderRecoveryLogic scans orders stuck in claimed status and re-queues them for processing.
|
|
type StuckOrderRecoveryLogic struct {
|
|
svc *svc.ServiceContext
|
|
}
|
|
|
|
func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogic {
|
|
return &StuckOrderRecoveryLogic{svc: svc}
|
|
}
|
|
|
|
// ProcessTask scans for orders stuck in claimed status for over 10 minutes,
|
|
// resets them to paid, and re-enqueues the activate task so they are retried
|
|
// independently of asynq's original retry counter (which may be exhausted).
|
|
func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
|
|
cutoff := time.Now().Add(-10 * time.Minute)
|
|
|
|
var stuckOrders []order.Order
|
|
if err := l.svc.DB.WithContext(ctx).
|
|
Model(&order.Order{}).
|
|
Where("status = ? AND updated_at < ?", OrderStatusClaimed, cutoff).
|
|
Find(&stuckOrders).Error; err != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to query stuck orders",
|
|
logger.Field("error", err.Error()),
|
|
)
|
|
return err
|
|
}
|
|
|
|
if len(stuckOrders) == 0 {
|
|
return nil
|
|
}
|
|
|
|
orderNos := make([]string, 0, len(stuckOrders))
|
|
for i := range stuckOrders {
|
|
orderNos = append(orderNos, stuckOrders[i].OrderNo)
|
|
}
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Found stuck claimed orders, recovering",
|
|
logger.Field("count", len(stuckOrders)),
|
|
logger.Field("order_nos", orderNos),
|
|
)
|
|
|
|
for i := range stuckOrders {
|
|
o := &stuckOrders[i]
|
|
|
|
// 终态守卫:OrderStatusClaimed(6) 与 orderStatusRefunded(6) 共用同一枚举值,
|
|
// 若该订单已写入 333 退款佣金日志,说明状态 6 表示「已退款」而非「短暂 claim」,
|
|
// 必须跳过,否则会把已退款订单重置为 5 + 重新入队 activate,导致重复退款。
|
|
// 详见 HIF-131 / HIF-132。
|
|
refunded, err := logmodel.HasRefundCommissionLog(l.svc.DB.WithContext(ctx), o.OrderNo)
|
|
if err != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to check refund log",
|
|
logger.Field("order_no", o.OrderNo),
|
|
logger.Field("error", err.Error()),
|
|
)
|
|
continue
|
|
}
|
|
if refunded {
|
|
logger.WithContext(ctx).Info("[StuckOrderRecovery] Skip refunded order (status=6 + refund log)",
|
|
logger.Field("order_no", o.OrderNo),
|
|
)
|
|
continue
|
|
}
|
|
|
|
result := l.svc.DB.WithContext(ctx).
|
|
Model(&order.Order{}).
|
|
Where("order_no = ? AND status = ?", o.OrderNo, OrderStatusClaimed).
|
|
Update("status", OrderStatusPaid)
|
|
if result.Error != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to reset order status",
|
|
logger.Field("order_no", o.OrderNo),
|
|
logger.Field("error", result.Error.Error()),
|
|
)
|
|
continue
|
|
}
|
|
if result.RowsAffected == 0 {
|
|
// Another process already handled this order
|
|
continue
|
|
}
|
|
|
|
// Invalidate order cache
|
|
o.Status = OrderStatusPaid
|
|
if err := l.svc.OrderModel.Update(ctx, o); err != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to update order cache",
|
|
logger.Field("order_no", o.OrderNo),
|
|
logger.Field("error", err.Error()),
|
|
)
|
|
}
|
|
|
|
// Re-enqueue activate task so the order gets processed regardless of asynq retry state
|
|
payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: o.OrderNo})
|
|
if err != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal task payload",
|
|
logger.Field("order_no", o.OrderNo),
|
|
logger.Field("error", err.Error()),
|
|
)
|
|
continue
|
|
}
|
|
if _, err = l.svc.Queue.EnqueueContext(ctx, asynq.NewTask(queueTypes.ForthwithActivateOrder, payload)); err != nil {
|
|
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue activate task",
|
|
logger.Field("order_no", o.OrderNo),
|
|
logger.Field("error", err.Error()),
|
|
)
|
|
} else {
|
|
logger.WithContext(ctx).Info("[StuckOrderRecovery] Re-enqueued activate task",
|
|
logger.Field("order_no", o.OrderNo),
|
|
)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|