5b9f384f81
P01:refundOrderLogic.RefundOrder 在事务内 FOR UPDATE 后、lockCommissionSource 前新增 333 退款日志扫描,命中即返回 OrderAlreadyRefunded(61006),不再写日志/扣 commission/ 改 order.status。 P02:堵住已退款订单状态被回退入口 - queue/logic/order/stuckOrderRecoveryLogic.go:批扫 status=6 时新增 333 日志守卫, 已退款订单不再被重置为 5 + 重新入队 activate(HIF-131 trace 中订单 53647 被刷回 5 的真凶) - queue/logic/order/activateOrderLogic.go:releaseClaim 同步加守卫做防御性兜底 新增 internal/model/log/refund.go 共享 helper HasRefundCommissionLog: type=33 + content LIKE 走索引粗筛,再 JSON 反序列化确认 content.type==333 AND content.order_no==orderNo,防 LIKE 子串误判。 测试:单元测试覆盖正常退款 / 已有 333 日志拒绝 / 子串误判防御 / 脏 JSON 容错; sqlmock 严格断言命中后事务序列只含 BEGIN/SELECT order FOR UPDATE/SELECT system_logs/ROLLBACK,无任何 commission 写入。 不做:calculateCommission、status 枚举拆分、表结构变更、支付通道 notify、用户余额回补。 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
|
|
}
|