5 Commits

Author SHA1 Message Date
247d7d9f6e 新增接口 2026-08-18 16:15:46 +08:00
d256f6d176 合并七月迭代分支 2026-08-18 14:53:29 +08:00
7029104e5c 让迁移套餐恢复月流量重置调度
缺少 next_reset_at 时,轮询根据已有激活时间或到期时间与套餐天数推算下一重置点;已有值通过查询条件和条件更新双重保护,不会被覆盖。

Constraint: 兼容迁移套餐缺少 activated_at 与 next_reset_at 的历史数据
Rejected: 单次 SQL 人工回填 | 后续迁移数据仍可能再次遗漏
Confidence: high
Scope-risk: narrow
Directive: 保持 next_reset_at 非空记录不可覆盖
Not-tested: 按用户要求未运行测试
2026-08-05 14:33:16 +08:00
a0de08d789 避免套餐过期后排队权益永久失联
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m6s
线上保持现有纯 Asynq 架构,以公平孤儿扫描和提交后投递消除永久饥饿及事务可见性竞态。

Constraint: 线上保持现有纯 Asynq 架构,不引入 Outbox、迁移或新任务基础设施。

Rejected: 事务内投递或扩大扫描 LIMIT | 无法消除竞态和永久饥饿。

Confidence: high

Scope-risk: narrow

Directive: 后续分支整合时按目标分支的套餐接续架构独立处理,不混用本热修实现。

Tested: go build ./...(退出码 0);git diff --check;openspec validate fix-main-package-activation-starvation --strict。

Not-tested: 按用户要求未新增、修改或运行自动化测试;线上 SQL、查询计划和日志待部署后核验。
2026-08-03 09:58:05 +08:00
1efb665619 fix: 修正排队顺延套餐激活时错误按下单时间计算生效日期
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 10m57s
activatePendingUsage 被"前一个主套餐到期后顺延激活下一个待生效套餐"和
"等待实名认证后激活"两种场景共用,但其中 ExpiryBase=from_purchase 计时
基准分支(REALNAME-04)本来只为后者设计,却被无差别套用到前者。

导致主套餐配置为 from_purchase 且需要排队等待前一个套餐到期才能生效的
套餐,激活时错误地把生效时间算成下单时间,而不是真正开始生效的那一刻,
使到期时间提前了排队等待的天数,客户少享受了相应天数的服务。

现改为只有当 usage.PendingRealnameActivation 为 true(确实是在等实名)
时才按 ExpiryBase 选择计时基准,纯排队顺延场景一律使用当前时刻,即顺延
语义。
2026-07-20 11:56:25 +09:00
27 changed files with 1781 additions and 16246 deletions

File diff suppressed because it is too large Load Diff

1045
README.md

File diff suppressed because it is too large Load Diff

View File

@@ -0,0 +1,44 @@
# Main 分支套餐接续饥饿热修
## 修复边界
本热修仅适用于 `main` 的纯 Asynq 套餐接续链路,不引入 Outbox、数据库迁移、新任务类型或新依赖。
- 孤儿扫描先在 PostgreSQL 中按卡或设备选择队首套餐,并排除仍有 `status IN (1,2)` 占位主套餐的载体,最后取 100 个真实孤儿。
- 旧主套餐及加油包状态事务提交后,再投递现有 `package:queue:activation` 任务。
- Redis 激活锁冲突返回套餐激活冲突错误,由现有 `MaxRetry(3)` 重试。
- 只有套餐实际从待生效推进为生效中时Handler 才记录“套餐激活成功”。
## 部署观察
部署 Worker 后至少观察两个套餐轮询周期:
1. `孤儿套餐扫描完成``orphan_count` 应能覆盖真实无占位套餐,不再固定被同一批占位载体挡住。
2. `已提交套餐激活任务` 后应出现实际激活、明确跳过或可重试错误,不再出现未改状态却打印成功。
3. Redis 锁冲突应进入 Asynq 重试,不应确认任务成功。
可使用以下只读 SQL 检查仍未恢复的真实孤儿数量:
```sql
SELECT COUNT(*) AS orphan_pending_count
FROM tb_package_usage AS pending
WHERE pending.status = 0
AND pending.master_usage_id IS NULL
AND pending.deleted_at IS NULL
AND (COALESCE(pending.iot_card_id, 0) > 0 OR COALESCE(pending.device_id, 0) > 0)
AND NOT EXISTS (
SELECT 1
FROM tb_package_usage AS occupied
WHERE occupied.status IN (1, 2)
AND occupied.master_usage_id IS NULL
AND occupied.deleted_at IS NULL
AND (
(COALESCE(pending.iot_card_id, 0) > 0 AND occupied.iot_card_id = pending.iot_card_id)
OR (COALESCE(pending.iot_card_id, 0) = 0 AND pending.device_id > 0 AND occupied.device_id = pending.device_id)
)
);
```
## 回滚
本次无数据库迁移。回滚热修提交并重新部署 Worker 即可;已经正确激活的套餐属于有效业务事实,不执行反向 SQL。

View File

@@ -8,6 +8,7 @@ import (
"github.com/bytedance/sonic"
"gorm.io/gorm"
"gorm.io/gorm/clause"
approvalapp "github.com/break/junhong_cmp_fiber/internal/application/approval"
"github.com/break/junhong_cmp_fiber/internal/model"
@@ -46,6 +47,94 @@ func NewOfflineCreationService(db *gorm.DB, approval approvalapp.Port, audit Rec
return &OfflineCreationService{db: db, approval: approval, audit: audit}
}
// TriggerHistorical 为历史待审批线下代充值补发一次企业微信审批。
func (s *OfflineCreationService) TriggerHistorical(ctx context.Context, recordID uint) (*CreateOfflineResult, error) {
if s == nil || s.db == nil || s.approval == nil || s.audit == nil || recordID == 0 {
return nil, errors.New(errors.CodeServiceUnavailable, "员工线下代充值审批能力未配置")
}
var record model.AgentRechargeRecord
if err := s.db.WithContext(ctx).First(&record, recordID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, errors.New(errors.CodeNotFound, "充值记录不存在")
}
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询历史线下代充值申请失败")
}
if record.PaymentMethod != constants.RechargeMethodOffline || record.Status != constants.RechargeStatusPending || record.ApprovalInstanceID != nil {
return nil, errors.New(errors.CodeConflict, "充值申请状态不允许补发审批")
}
account, shop, wallet, err := s.loadHistoricalFacts(ctx, &record)
if err != nil {
return nil, err
}
preparation, err := s.approval.Prepare(ctx, approvalapp.PrepareRequest{
BusinessType: constants.ApprovalBusinessTypeOfflineRecharge, SubmitterAccountID: record.UserID,
CorrelationID: record.RechargeNo,
})
if err != nil {
return nil, err
}
command := CreateOfflineCommand{
SubmitterAccountID: record.UserID, SubmitterUserType: account.UserType, ShopID: record.ShopID,
RechargeNo: record.RechargeNo, Amount: record.Amount,
PaymentVoucherKeys: []string(record.PaymentVoucherKey), Remark: record.Remark,
}
submitterSnapshot, requestSnapshot, err := offlineApprovalSnapshots(command, account.Username, shop.ShopName)
if err != nil {
return nil, err
}
var approvalStatus int
err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var current model.AgentRechargeRecord
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).First(&current, recordID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return errors.New(errors.CodeNotFound, "充值记录不存在")
}
return errors.Wrap(errors.CodeDatabaseError, err, "锁定历史线下代充值申请失败")
}
if current.PaymentMethod != constants.RechargeMethodOffline || current.Status != constants.RechargeStatusPending || current.ApprovalInstanceID != nil {
return errors.New(errors.CodeConflict, "充值申请状态不允许补发审批")
}
reference, err := s.approval.CreateInTx(ctx, tx, approvalapp.CreateRequest{
Preparation: preparation, BusinessType: constants.ApprovalBusinessTypeOfflineRecharge,
BusinessID: current.ID, SubmitterAccountID: current.UserID,
SubmitterSnapshot: submitterSnapshot, RequestSnapshot: requestSnapshot,
CorrelationID: current.RechargeNo,
})
if err != nil {
return err
}
result := tx.WithContext(ctx).Model(&model.AgentRechargeRecord{}).
Where("id = ? AND payment_method = ? AND status = ? AND approval_instance_id IS NULL", current.ID, constants.RechargeMethodOffline, constants.RechargeStatusPending).
Update("approval_instance_id", reference.InstanceID)
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "关联线下代充值审批实例失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "线下代充值审批实例关联已变化")
}
current.ApprovalInstanceID = &reference.InstanceID
record = current
approvalStatus = reference.Status
var instance model.ApprovalInstance
if err := tx.WithContext(ctx).First(&instance, reference.InstanceID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询线下代充值审批审计快照失败")
}
return s.audit.WriteAgentRecharge(ctx, tx, RechargeAudit{
ActionCode: constants.AuditActionAgentRechargeCreated, Summary: "补发员工线下代充值审批",
Record: &current, Approval: &instance, Wallet: wallet,
AfterData: map[string]any{"status": current.Status, "approval_instance_id": current.ApprovalInstanceID},
})
})
if err != nil {
return nil, err
}
return &CreateOfflineResult{
Record: &record, ShopName: shop.ShopName, SubmitterName: account.Username, ApprovalStatus: approvalStatus,
}, nil
}
// Execute 在业务写入前校验审批渠道,并在同一事务保存充值申请、审批实例和提交 Outbox。
func (s *OfflineCreationService) Execute(ctx context.Context, command CreateOfflineCommand) (*CreateOfflineResult, error) {
if s == nil || s.db == nil || s.approval == nil || s.audit == nil {
@@ -140,6 +229,38 @@ func validateCreateOfflineCommand(command CreateOfflineCommand) error {
return nil
}
func (s *OfflineCreationService) loadHistoricalFacts(
ctx context.Context, record *model.AgentRechargeRecord,
) (*model.Account, *model.Shop, *model.AgentWallet, error) {
if record == nil || record.UserID == 0 || record.ShopID == 0 || record.AgentWalletID == 0 {
return nil, nil, nil, errors.New(errors.CodeInvalidParam)
}
var account model.Account
if err := s.db.WithContext(ctx).Where("id = ? AND status = ?", record.UserID, constants.StatusEnabled).First(&account).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, nil, nil, errors.New(errors.CodeForbidden, "原创建账号不可用")
}
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询历史线下代充值创建人失败")
}
var shop model.Shop
if err := s.db.WithContext(ctx).Where("id = ?", record.ShopID).First(&shop).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, nil, nil, errors.New(errors.CodeNotFound, "目标店铺不存在")
}
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询历史线下代充值目标店铺失败")
}
var wallet model.AgentWallet
if err := s.db.WithContext(ctx).
Where("id = ? AND shop_id = ? AND wallet_type = ? AND status = ?", record.AgentWalletID, record.ShopID, constants.AgentWalletTypeMain, constants.AgentWalletStatusNormal).
First(&wallet).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, nil, nil, errors.New(errors.CodeWalletNotFound, "原充值主钱包不存在或不可用")
}
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询历史线下代充值主钱包失败")
}
return &account, &shop, &wallet, nil
}
func (s *OfflineCreationService) loadCreationFacts(
ctx context.Context,
command CreateOfflineCommand,

View File

@@ -8,6 +8,7 @@ import (
"github.com/bytedance/sonic"
"gorm.io/gorm"
"gorm.io/gorm/clause"
approvalapp "github.com/break/junhong_cmp_fiber/internal/application/approval"
"github.com/break/junhong_cmp_fiber/internal/model"
@@ -55,6 +56,99 @@ func NewCreationService(db *gorm.DB, approval approvalapp.Port, audit AuditWrite
}
// Execute 在业务写入前校验审批渠道,并在同一事务冻结退款事实和审批事实。
// TriggerHistorical 为历史待审批退款补发一次企业微信审批。
func (s *CreationService) TriggerHistorical(ctx context.Context, refundID uint) (*CreateResult, error) {
if s == nil || s.db == nil || s.approval == nil || s.audit == nil || refundID == 0 {
return nil, errors.New(errors.CodeServiceUnavailable, "退款审批能力未配置")
}
var refund model.RefundRequest
if err := s.db.WithContext(ctx).First(&refund, refundID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, errors.New(errors.CodeNotFound, "退款申请不存在")
}
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询历史退款申请失败")
}
if refund.Status != model.RefundStatusPending || refund.ApprovalInstanceID != nil {
return nil, errors.New(errors.CodeConflict, "退款申请状态不允许补发审批")
}
account, err := s.loadSubmitter(ctx, refund.Creator)
if err != nil {
return nil, err
}
var order model.Order
if err := s.db.WithContext(ctx).First(&order, refund.OrderID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, errors.New(errors.CodeNotFound, "退款关联订单不存在")
}
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询退款关联订单失败")
}
preparation, err := s.approval.Prepare(ctx, approvalapp.PrepareRequest{
BusinessType: constants.ApprovalBusinessTypeRefund, SubmitterAccountID: refund.Creator,
CorrelationID: refund.RefundNo,
})
if err != nil {
return nil, err
}
submitterSnapshot, requestSnapshot, err := refundSnapshots(&refund, account)
if err != nil {
return nil, err
}
var approvalStatus int
err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var current model.RefundRequest
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).First(&current, refundID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return errors.New(errors.CodeNotFound, "退款申请不存在")
}
return errors.Wrap(errors.CodeDatabaseError, err, "锁定历史退款申请失败")
}
if current.Status != model.RefundStatusPending || current.ApprovalInstanceID != nil {
return errors.New(errors.CodeConflict, "退款申请状态不允许补发审批")
}
var currentOrder model.Order
if err := tx.WithContext(ctx).First(&currentOrder, current.OrderID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询退款关联订单失败")
}
reference, err := s.approval.CreateInTx(ctx, tx, approvalapp.CreateRequest{
Preparation: preparation, BusinessType: constants.ApprovalBusinessTypeRefund,
BusinessID: current.ID, SubmitterAccountID: current.Creator,
SubmitterSnapshot: submitterSnapshot, RequestSnapshot: requestSnapshot,
CorrelationID: current.RefundNo,
})
if err != nil {
return err
}
result := tx.WithContext(ctx).Model(&model.RefundRequest{}).
Where("id = ? AND status = ? AND approval_instance_id IS NULL", current.ID, model.RefundStatusPending).
Update("approval_instance_id", reference.InstanceID)
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "关联退款审批实例失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "退款审批实例关联已变化")
}
current.ApprovalInstanceID = &reference.InstanceID
refund = current
order = currentOrder
approvalStatus = reference.Status
var instance model.ApprovalInstance
if err := tx.WithContext(ctx).First(&instance, reference.InstanceID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询退款审批审计快照失败")
}
return s.audit.WriteRefundApplication(ctx, tx, ApplicationAudit{
Refund: &current, Order: &currentOrder, Approval: &instance, Submitter: account,
})
})
if err != nil {
return nil, err
}
return &CreateResult{Refund: &refund, SubmitterName: account.Username, ApprovalStatus: approvalStatus}, nil
}
func (s *CreationService) Execute(ctx context.Context, command CreateCommand) (*CreateResult, error) {
if s == nil || s.db == nil || s.approval == nil || s.audit == nil {
return nil, errors.New(errors.CodeServiceUnavailable, "退款审批能力未配置")

View File

@@ -143,6 +143,20 @@ func (h *AgentRechargeHandler) Get(c *fiber.Ctx) error {
return response.Success(c, result)
}
// TriggerApproval 主动补发历史线下代理充值审批。
// POST /api/admin/agent-recharges/:id/trigger-approval
func (h *AgentRechargeHandler) TriggerApproval(c *fiber.Ctx) error {
id, err := strconv.ParseUint(c.Params("id"), 10, 64)
if err != nil || id == 0 {
return errors.New(errors.CodeInvalidParam, "无效的充值记录ID")
}
result, err := h.service.TriggerApproval(c.UserContext(), uint(id))
if err != nil {
return err
}
return response.Success(c, result)
}
// PaymentStatus 查询代理充值本地支付与到账状态。
// GET /api/admin/agent-recharges/:id/payment-status
func (h *AgentRechargeHandler) PaymentStatus(c *fiber.Ctx) error {

View File

@@ -69,6 +69,20 @@ func (h *RefundHandler) GetByID(c *fiber.Ctx) error {
return response.Success(c, result)
}
// TriggerApproval 主动补发历史退款审批
// POST /api/admin/refunds/:id/trigger-approval
func (h *RefundHandler) TriggerApproval(c *fiber.Ctx) error {
id, err := strconv.ParseUint(c.Params("id"), 10, 64)
if err != nil || id == 0 {
return errors.New(errors.CodeInvalidParam, "无效的退款申请ID")
}
result, err := h.service.TriggerApproval(c.UserContext(), uint(id))
if err != nil {
return err
}
return response.Success(c, result)
}
// Approve 审批通过退款申请
// POST /api/admin/refunds/:id/approve
func (h *RefundHandler) Approve(c *fiber.Ctx) error {

View File

@@ -1,82 +0,0 @@
package audit
import (
"context"
"encoding/json"
"testing"
"gorm.io/gorm"
accessauditapp "github.com/break/junhong_cmp_fiber/internal/application/accessaudit"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/pkg/auditfailure"
"github.com/break/junhong_cmp_fiber/pkg/constants"
)
func TestAppendFailureDoesNotReturnToBusiness(t *testing.T) {
writer := NewWriter(nil, nil)
input := AppendInput{ActionCode: "missing_action"}
before := auditfailure.SecondaryWriteFailureCount()
if err := writer.Append(context.Background(), nil, input); err != nil {
t.Fatalf("Append 返回审计失败: %v", err)
}
if err := writer.WriteAccessChange(context.Background(), nil, accessauditapp.ChangeAudit{
ActionCode: constants.AuditActionPersonalCustomerAssetBound,
OperatorID: 1,
}); err != nil {
t.Fatalf("资源构造失败返回业务: %v", err)
}
if got := auditfailure.SecondaryWriteFailureCount(); got != before+2 {
t.Fatalf("二次失败记录次数 = %d, want %d", got, before+2)
}
if _, err := writer.AppendAndGet(context.Background(), nil, input); err == nil {
t.Fatal("AppendAndGet 未保留错误语义")
}
}
func TestPersonalCustomerAssetBoundProjectsOnlyPersonalResources(t *testing.T) {
action, ok := NewRegistry().Action(constants.AuditActionPersonalCustomerAssetBound)
if !ok {
t.Fatal("未注册个人客户资产绑定审计动作")
}
resources, err := accessResources(accessauditapp.ChangeAudit{
ActionCode: constants.AuditActionPersonalCustomerAssetBound,
PersonalCustomer: &model.PersonalCustomer{Model: gorm.Model{ID: 1}, Nickname: "客户"},
PersonalDevices: []accessauditapp.PersonalCustomerDeviceChange{{
Binding: &model.PersonalCustomerDevice{Model: gorm.Model{ID: 2}, CustomerID: 1, VirtualNo: "DEVICE-1"},
}},
PersonalICCIDs: []accessauditapp.PersonalCustomerICCIDChange{{
Binding: &model.PersonalCustomerICCID{Model: gorm.Model{ID: 3}, CustomerID: 1, ICCID: "ICCID-1"},
}},
SubjectVisibility: constants.AuditSubjectDetail,
SubjectSummary: "绑定个人客户资产",
SubjectData: map[string]any{"asset_type": constants.AuditResourceIotCard, "asset_id": uint(9)},
}, action.PrimaryResource)
if err != nil {
t.Fatalf("构造绑定审计资源失败: %v", err)
}
projected, err := NewWriter(nil, nil).buildResources(resources, action)
if err != nil {
t.Fatalf("构造绑定审计投影失败: %v", err)
}
want := map[string]bool{
constants.AuditResourcePersonalCustomer: true,
constants.AuditResourcePersonalCustomerDevice: true,
constants.AuditResourcePersonalCustomerICCID: true,
}
for _, resource := range projected {
if resource.ResourceType == constants.AuditResourceIotCard || resource.ResourceType == constants.AuditResourceDevice {
t.Fatalf("绑定审计投影包含内部资源: %s", resource.ResourceType)
}
delete(want, resource.ResourceType)
if resource.ResourceType == constants.AuditResourcePersonalCustomer {
var subjectData map[string]any
if err := json.Unmarshal(resource.SubjectData, &subjectData); err != nil || resource.SubjectVisibility != constants.AuditSubjectDetail || subjectData["asset_type"] != constants.AuditResourceIotCard || subjectData["asset_id"] != float64(9) {
t.Fatalf("主个人客户主体投影不完整: %#v", resource)
}
}
}
for resourceType := range want {
t.Fatalf("绑定审计投影缺少合法资源: %s", resourceType)
}
}

View File

@@ -223,29 +223,14 @@ func (h *PackageActivationHandler) findAndActivateOrphanPackages(ctx context.Con
count := 0
for _, usage := range orphanUsages {
carrierType, carrierID := h.getCarrierInfo(usage)
activated, activationErr := h.activationService.ActivateNextPendingMainPackage(ctx, carrierType, carrierID)
if activationErr != nil {
h.logger.Warn("孤儿套餐同步激活失败",
if err := h.enqueueActivationTask(ctx, usage.ID, carrierType, carrierID, "orphan_recovery"); err != nil {
h.logger.Warn("提交孤儿套餐激活任务失败",
zap.Uint("package_usage_id", usage.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "orphan_recovery"),
zap.Error(activationErr))
zap.Error(err))
continue
}
if !activated {
h.logger.Info("孤儿套餐本轮未激活",
zap.Uint("package_usage_id", usage.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "orphan_recovery"))
continue
}
h.logger.Info("孤儿套餐同步激活成功",
zap.Uint("package_usage_id", usage.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "orphan_recovery"))
count++
}
@@ -268,7 +253,7 @@ func (h *PackageActivationHandler) findExpiredMainPackages(ctx context.Context)
}
// processExpiredPackage 处理单个过期套餐
// 流程:先同步最新流量 → 事务内标记过期和失效加油包 → 提交后同步接续并触发停机检查
// 流程:先同步最新流量 → 事务内标记过期和失效加油包 → 提交后投递下一套餐并触发停机检查
func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pkg *model.PackageUsage) error {
carrierType, carrierID := h.getCarrierInfo(pkg)
@@ -324,28 +309,9 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
return nil
}
// 事务提交后再投递,确保消费者只能读取到旧套餐已经过期的状态。
if carrierType != "" && carrierID > 0 {
activated, activationErr := h.activationService.ActivateNextPendingMainPackage(ctx, carrierType, carrierID)
if activationErr != nil {
h.logger.Warn("过期后同步接续套餐失败",
zap.Uint("expired_package_usage_id", pkg.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "expired_package"),
zap.Error(activationErr))
} else if activated {
h.logger.Info("过期后同步接续套餐成功",
zap.Uint("expired_package_usage_id", pkg.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "expired_package"))
} else {
h.logger.Info("过期后本轮未接续套餐",
zap.Uint("expired_package_usage_id", pkg.ID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.String("activation_source", "expired_package"))
}
activationErr := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
h.triggerStopAfterExpiry(ctx, carrierType, carrierID)
return activationErr
}
@@ -401,6 +367,30 @@ func (h *PackageActivationHandler) getCarrierInfo(pkg *model.PackageUsage) (stri
return "", 0
}
// activateNextPackage 提交下一个待生效主套餐的异步激活任务。
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) error {
var nextPkg model.PackageUsage
query := tx.Where("status = ?", constants.PackageUsageStatusPending).
Where("master_usage_id IS NULL").
Order("priority ASC, created_at ASC, id ASC").
Limit(1)
if carrierType == constants.AssetTypeIotCard {
query = query.Where("iot_card_id = ?", carrierID)
} else if carrierType == constants.AssetTypeDevice {
query = query.Where("device_id = ?", carrierID)
}
if err := query.First(&nextPkg).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil
}
return err
}
return h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue")
}
// triggerStopAfterExpiry 套餐过期后异步触发停机检查
// 仅在确认无后续生效套餐时有效CheckAndStopCard 内部有幂等保护,重复调用安全
func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, carrierType string, carrierID uint) {

View File

@@ -61,6 +61,14 @@ func registerAgentRechargeRoutes(router fiber.Router, handler *admin.AgentRechar
Auth: true,
})
Register(group, doc, groupPath, "POST", "/:id/trigger-approval", handler.TriggerApproval, RouteSpec{
Summary: "补发历史线下代理充值审批",
Tags: []string{"代理预充值"},
Input: new(dto.IDReq),
Output: new(dto.AgentRechargeResponse),
Auth: true,
})
Register(group, doc, groupPath, "POST", "/:id/offline-pay", handler.OfflinePay, RouteSpec{
Summary: "确认线下充值",
Tags: []string{"代理预充值"},

View File

@@ -49,6 +49,14 @@ func registerRefundRoutes(router fiber.Router, handler *admin.RefundHandler, doc
Auth: true,
})
Register(refund, doc, groupPath, "POST", "/:id/trigger-approval", handler.TriggerApproval, RouteSpec{
Summary: "补发历史退款审批",
Tags: []string{"退款管理"},
Input: new(dto.RefundIDRequest),
Output: new(dto.RefundResponse),
Auth: true,
})
Register(refund, doc, groupPath, "POST", "/:id/approve", handler.Approve, RouteSpec{
Summary: "审批通过退款申请",
Tags: []string{"退款管理"},

View File

@@ -434,6 +434,27 @@ func (s *Service) appendCreditedAudit(ctx context.Context, tx *gorm.DB, record *
})
}
// TriggerApproval 为历史线下代理充值主动补发企业微信审批。
func (s *Service) TriggerApproval(ctx context.Context, id uint) (*dto.AgentRechargeResponse, error) {
if s.offlineCreation == nil {
return nil, errors.New(errors.CodeServiceUnavailable, "员工线下代充值审批能力未配置")
}
record, err := s.agentRechargeStore.GetByID(ctx, id)
if err != nil {
return nil, errors.New(errors.CodeNotFound, "充值记录不存在")
}
result, err := s.offlineCreation.TriggerHistorical(ctx, record.ID)
if err != nil {
return nil, err
}
resp := toResponse(result.Record, result.ShopName)
resp.SubmitterName = result.SubmitterName
resp.ApprovalProvider = constants.IntegrationProviderWeCom
resp.ApprovalStatus = &result.ApprovalStatus
resp.ApprovalStatusName = constants.GetApprovalStatusName(result.ApprovalStatus)
return resp, nil
}
// GetByID 根据ID查询充值订单详情
// GET /api/admin/agent-recharges/:id
func (s *Service) GetByID(ctx context.Context, id uint) (*dto.AgentRechargeResponse, error) {

View File

@@ -1,83 +0,0 @@
package customer_binding
import (
"context"
"database/sql"
"database/sql/driver"
"io"
"testing"
"gorm.io/driver/postgres"
"gorm.io/gorm"
accessauditapp "github.com/break/junhong_cmp_fiber/internal/application/accessaudit"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/pkg/constants"
)
func init() { sql.Register("customer_binding_audit_test", customerAuditDriver{}) }
type customerAuditDriver struct{}
func (customerAuditDriver) Open(string) (driver.Conn, error) { return customerAuditConn{}, nil }
type customerAuditConn struct{}
func (customerAuditConn) Prepare(string) (driver.Stmt, error) { return nil, driver.ErrSkip }
func (customerAuditConn) Close() error { return nil }
func (customerAuditConn) Begin() (driver.Tx, error) { return nil, driver.ErrSkip }
func (customerAuditConn) QueryContext(context.Context, string, []driver.NamedValue) (driver.Rows, error) {
return &customerAuditRows{}, nil
}
type customerAuditRows struct{ sent bool }
func (*customerAuditRows) Columns() []string { return []string{"id", "nickname"} }
func (r *customerAuditRows) Close() error { return nil }
func (r *customerAuditRows) Next(dest []driver.Value) error {
if r.sent {
return io.EOF
}
r.sent = true
dest[0], dest[1] = int64(7), "客户"
return nil
}
type captureAuditWriter struct{ change accessauditapp.ChangeAudit }
func (w *captureAuditWriter) WriteAccessChange(_ context.Context, _ *gorm.DB, change accessauditapp.ChangeAudit) error {
w.change = change
return nil
}
func TestWriteBindingAuditOmitsInternalAssets(t *testing.T) {
db, err := sql.Open("customer_binding_audit_test", "")
if err != nil {
t.Fatal(err)
}
defer db.Close()
tx, err := gorm.Open(postgres.New(postgres.Config{Conn: db}), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
writer := &captureAuditWriter{}
service := &Service{accessAudit: writer}
personalDevices := []accessauditapp.PersonalCustomerDeviceChange{{Binding: &model.PersonalCustomerDevice{Model: gorm.Model{ID: 2}, CustomerID: 7, VirtualNo: "DEVICE-1"}}}
personalICCIDs := []accessauditapp.PersonalCustomerICCIDChange{{Binding: &model.PersonalCustomerICCID{Model: gorm.Model{ID: 3}, CustomerID: 7, ICCID: "ICCID-1"}}}
cards := []accessauditapp.IotCardChange{{Card: &model.IotCard{Model: gorm.Model{ID: 9}}}}
devices := []accessauditapp.DeviceChange{{Device: &model.Device{Model: gorm.Model{ID: 10}}}}
if err := service.writeBindingAudit(context.Background(), tx, constants.AuditActionPersonalCustomerAssetBound, "绑定个人客户资产", 7, personalDevices, personalICCIDs, cards, devices); err != nil {
t.Fatalf("写入绑定审计失败: %v", err)
}
change := writer.change
if len(change.Cards) != 0 || len(change.Devices) != 0 {
t.Fatalf("绑定审计泄露内部资源: Cards=%d Devices=%d", len(change.Cards), len(change.Devices))
}
if change.PersonalCustomer == nil || change.PersonalCustomer.ID != 7 || len(change.PersonalDevices) != 1 || len(change.PersonalICCIDs) != 1 {
t.Fatalf("绑定审计未保留个人客户字段: %#v", change)
}
if change.SubjectData["asset_type"] != constants.AuditResourceIotCard || change.SubjectData["asset_id"] != uint(9) {
t.Fatalf("绑定审计未保留主体摘要: %#v", change.SubjectData)
}
}

View File

@@ -238,6 +238,27 @@ func (s *Service) List(ctx context.Context, req *dto.RefundListRequest) (*dto.Re
}
// GetByID 根据 ID 查询退款申请详情
// TriggerApproval 为历史退款申请主动补发企业微信审批。
func (s *Service) TriggerApproval(ctx context.Context, id uint) (*dto.RefundResponse, error) {
if s.refundApprovalCreation == nil {
return nil, errors.New(errors.CodeServiceUnavailable, "退款审批能力未配置")
}
refund, err := s.refundStore.GetByIDForOperation(ctx, id)
if err != nil {
return nil, errors.New(errors.CodeNotFound, "退款申请不存在")
}
result, err := s.refundApprovalCreation.TriggerHistorical(ctx, refund.ID)
if err != nil {
return nil, err
}
resp := buildRefundResponse(result.Refund)
resp.SubmitterName = result.SubmitterName
resp.ApprovalProvider = constants.IntegrationProviderWeCom
resp.ApprovalStatus = &result.ApprovalStatus
resp.ApprovalStatusName = constants.GetApprovalStatusName(result.ApprovalStatus)
return resp, nil
}
func (s *Service) GetByID(ctx context.Context, id uint) (*dto.RefundResponse, error) {
refund, err := s.refundStore.GetByID(ctx, id)
if err != nil {

View File

@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-08-18

View File

@@ -0,0 +1,45 @@
## Context
历史记录在 `tb_refund_request``tb_agent_recharge_record` 中保持待审批但 `approval_instance_id` 为空。现有新建用例已经在单个事务中创建通用审批实例、企业微信上下文、审批提交 Outbox并回填该关联两个业务表的 `approval_instance_id` 均有唯一索引,通用审批实例还以业务类型和业务 ID 唯一。
## Goals / Non-Goals
**Goals:**
- 以最小增量复用现有审批创建和可靠提交链路补发历史记录。
- 以原创建账号构造企业微信发起人及审批快照。
- 使并发请求和重复请求均不会形成第二张审批单。
**Non-Goals:**
- 不批量扫描或自动补发历史记录。
- 不改变既有审批终态、企业微信提交重试或人工审批接口。
- 不新增迁移、重置既有审批关联,或为提交失败创建第二张审批单。
## Decisions
### 在各业务审批创建用例中增加历史记录发起入口
退款和线下代理充值分别新增面向既有记录的 Application 用例入口,复用各自已有的快照构造、提交人校验、通用审批 `Prepare`/`CreateInTx`、审计及 DTO 组装逻辑。Handler 只解析路径 ID 并调用服务Service 负责加载完整业务事实和调用 Application。
选择按业务保留两个小入口,而不引入跨退款/充值的通用“历史审批补发器”:二者的资格条件、快照和关联事实不同,现有两个创建用例已是最短复用边界。
### 以事务内条件更新和既有唯一约束保证一次性
发起前可在事务外执行审批渠道预检;事务内必须重新读取或条件更新业务记录,要求 `status=待审批 AND approval_instance_id IS NULL`,再创建通用审批及渠道上下文/Outbox并回填 `approval_instance_id`。任一环节失败回滚,不消耗发起资格;成功提交后,由业务表关联唯一索引和通用审批业务唯一索引共同拒绝并发的第二次创建。
不增加“已尝试”字段:用户确认以成功创建审批实例作为一次性边界,已有唯一关联就是持久化且可恢复的事实源。
### 发起人和授权语义
企业微信发起人固定为业务记录 `Creator`,不使用点击接口的账号;该账号不可用时失败关闭。接口沿用各自当前路由组的账号类型授权,不扩大既有退款或代理充值管理入口的访问范围。返回值沿用现有详情 DTO 的审批摘要字段,避免新增响应类型。
## Risks / Trade-offs
- [原创建账号已禁用或未绑定企业微信] → 不创建任何审批事实并返回错误;维护者修复账号/绑定后可再次操作。
- [两个请求同时发起] → 事务条件和数据库唯一约束确保仅一个提交成功,调用方对另一个请求按冲突处理。
- [提交 Outbox 后企微调用结果未知] → 沿用已有结果未知恢复流程,禁止通过本接口重建审批。
## Migration Plan
1. 发布 API 与 Worker 均包含该版本的应用代码,确保 Outbox 消费者已注册。
2. 维护者在生产环境按发布运行说明,通过列表筛选待审批历史记录后逐单调用新接口,并核对返回的审批摘要与审计/Outbox 事实。
3. 如需回滚,仅停止暴露新路由并回滚应用二进制;已成功创建的审批实例继续由既有 Worker 流程处理,不删除审批关联或重新发起。

View File

@@ -0,0 +1,28 @@
## Why
七月迭代上线前已创建且仍待审批的退款申请、员工线下代充值申请未关联通用审批实例,无法进入企业微信审批流。需要由管理员按单主动补发,同时避免同一业务重复创建审批单。
## What Changes
- 为待审批且尚未关联审批实例的历史退款申请新增主动发起企业微信审批接口。
- 为待审批、线下支付且尚未关联审批实例的历史代理充值申请新增主动发起企业微信审批接口。
- 主动发起时复用原业务创建人作为企业微信审批发起人;原创建人不可用或审批场景不可用时不创建审批实例。
- 在同一事务创建通用审批实例、企业微信上下文、提交 Outbox 并回填业务记录的 `approval_instance_id`,以该唯一关联保证成功创建后不可再次发起。
- 仅在业务保持待审批状态时允许主动发起;已关联审批实例、非线下充值或非待审批记录均拒绝。
## Capabilities
### New Capabilities
- 无。
### Modified Capabilities
- `order-refund-exchange`: 退款申请可对历史待审批且未关联审批实例的记录主动创建一次企业微信审批。
- `agent-funds-commission`: 历史待审批线下代理充值申请可主动创建一次企业微信审批。
## Impact
- 路由、退款与代理充值 Handler/Service以及审批创建 Application 用例。
- 新增两个后台 API 并同步 OpenAPI 文档生成入口。
- 复用现有通用审批、企业微信审批上下文、Outbox、审计和既有 `approval_instance_id` 唯一索引;不新增外部依赖或数据库表结构。

View File

@@ -0,0 +1,21 @@
## ADDED Requirements
### Requirement: 历史待审批线下代理充值可主动接入企业微信审批
系统 SHALL 提供 `POST /api/admin/agent-recharges/{id}/trigger-approval`,使具有既有代理充值管理访问权限的后台账号可为历史线下代理充值申请主动创建企业微信审批。系统 MUST 仅在线下充值记录处于待审批状态且 `approval_instance_id` 为空时创建审批;审批发起人 MUST 使用该充值记录的原创建账号。创建成功后,系统 MUST 原子保存唯一审批实例、审批提交请求及充值记录的审批实例关联,并返回更新后的充值申请审批摘要。
#### Scenario: 主动发起历史线下代理充值审批成功
- **GIVEN** 线下代理充值申请处于待审批状态、未关联审批实例,且其原创建账号和企业微信线下充值审批场景均可用
- **WHEN** 有既有代理充值管理访问权限的后台账号请求 `POST /api/admin/agent-recharges/{id}/trigger-approval`
- **THEN** 系统创建以原创建账号为发起人的唯一企业微信审批并返回审批摘要,后续由既有可靠提交流程提交至企业微信
#### Scenario: 在线、非待审批或已发起记录被拒绝
- **WHEN** 请求主动发起的充值记录不是线下充值、不是待审批状态或已关联审批实例
- **THEN** 系统返回状态冲突且不创建新的审批实例或提交请求
#### Scenario: 并发主动发起同一充值审批
- **WHEN** 两个请求同时为同一符合条件的线下代理充值申请主动发起审批
- **THEN** 系统至多创建一个审批实例和一个审批提交请求,未成功创建关联的请求返回冲突
#### Scenario: 原创建人或审批渠道不可用
- **WHEN** 充值申请原创建账号不可用,或企业微信线下充值审批场景不可用
- **THEN** 系统返回相应错误,充值申请保持未关联审批实例,修复条件后可再次发起

View File

@@ -0,0 +1,21 @@
## ADDED Requirements
### Requirement: 历史待审批退款可主动接入企业微信审批
系统 SHALL 提供 `POST /api/admin/refunds/{id}/trigger-approval`,使具有既有退款管理访问权限的后台账号可为历史退款申请主动创建企业微信审批。系统 MUST 仅在退款申请状态为待审批且 `approval_instance_id` 为空时创建审批;审批发起人 MUST 使用该退款申请的原创建账号。创建成功后,系统 MUST 原子保存唯一审批实例、审批提交请求及退款申请的审批实例关联,并返回更新后的退款申请审批摘要。
#### Scenario: 主动发起历史退款审批成功
- **GIVEN** 退款申请处于待审批状态、未关联审批实例,且其原创建账号和企业微信退款审批场景均可用
- **WHEN** 有既有退款管理访问权限的后台账号请求 `POST /api/admin/refunds/{id}/trigger-approval`
- **THEN** 系统创建以原创建账号为发起人的唯一企业微信审批并返回审批摘要,后续由既有可靠提交流程提交至企业微信
#### Scenario: 非待审批或已发起记录被拒绝
- **WHEN** 请求主动发起的退款申请不是待审批状态或已关联审批实例
- **THEN** 系统返回状态冲突且不创建新的审批实例或提交请求
#### Scenario: 并发主动发起同一退款审批
- **WHEN** 两个请求同时为同一符合条件的退款申请主动发起审批
- **THEN** 系统至多创建一个审批实例和一个审批提交请求,未成功创建关联的请求返回冲突
#### Scenario: 原创建人或审批渠道不可用
- **WHEN** 退款申请原创建账号不可用,或企业微信退款审批场景不可用
- **THEN** 系统返回相应错误,退款申请保持未关联审批实例,修复条件后可再次发起

View File

@@ -0,0 +1,16 @@
## 1. 审批补发用例
- [x] 1.1 在退款审批 Application 中实现历史待审批退款的主动发起:加载原创建人和订单事实、复用既有审批快照与 `Prepare`/`CreateInTx` 链路,并在同一事务内按待审批且未关联审批实例的条件回填关联和审计。
- [x] 1.2 在线下代理充值 Application 中实现历史待审批充值的主动发起校验线下支付、待审批和未关联审批实例加载原创建人、店铺和钱包事实并复用既有审批创建、快照、Outbox 与审计链路。
- [x] 1.3 在退款和代理充值 Service 中接入补发用例,复核既有路由权限与资源查询范围,向调用方返回包含审批摘要的既有 DTO将并发或已关联审批实例映射为状态冲突。
## 2. HTTP 入口与文档
- [x] 2.1 在退款 Handler 和路由注册 `POST /api/admin/refunds/{id}/trigger-approval`,完成路径 ID 绑定并交由 Service 处理。
- [x] 2.2 在代理充值 Handler 和路由注册 `POST /api/admin/agent-recharges/{id}/trigger-approval`,完成路径 ID 绑定并交由 Service 处理。
- [x] 2.3 同步 `cmd/api/docs.go``cmd/gendocs/main.go` 所依赖的路由元数据,确保两个接口及其响应模型生成到 OpenAPI 文档。
## 3. 验证
- [x] 3.1 以隔离环境或最小可运行验证覆盖:两个符合资格的历史记录各只创建一次审批实例和提交 Outbox非待审批、已关联、在线充值、原创建人/场景不可用及并发重复请求不创建第二实例。
- [ ] 3.2 执行 `gofmt -w`(变更的 Go 文件)、`go build ./cmd/api ./cmd/worker``go run cmd/gendocs/main.go``openspec validate --all``./scripts/context-health.sh`

View File

@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-08-03

View File

@@ -0,0 +1,106 @@
## Context
线上只读数据得到旧孤儿扫描窗口 `100/100/0`:固定取出的 100 条待生效记录全部仍有 `status IN (1,2)` 占位主套餐,而真正无占位套餐的记录位于窗口之外。当前实现先 `LIMIT 100`,再逐条查询占位状态,因此同一批无效候选会永久挡住真实孤儿。
过期路径还在数据库事务内调用 `enqueueActivationTask`。Asynq 消费者可能早于事务提交读取旧主套餐,随后 `ActivateSpecificPackage` 判断“已有生效主套餐”并返回 `nil`Handler 继续记录“套餐激活成功”任务不再重试。Redis 锁冲突也返回 `nil`,存在另一条假成功路径。
本热修保持当前线上 `Polling Handler → GORM transaction → Asynq → Package Activation Service` 架构,不引入新的可靠投递设施。
## Goals / Non-Goals
**Goals:**
- 每轮 100 个恢复名额只用于真实孤儿载体,同一载体只选择队首套餐。
- 旧套餐过期事实提交后才投递 Asynq消除事务可见性竞态。
- Redis 锁冲突返回错误,由现有 Asynq 重试。
- 成功日志只对应实际激活或明确的已完成幂等事实。
**Non-Goals:**
- 不新增 Outbox、迁移、索引、依赖、队列或任务类型。
- 不重构套餐购买、实名激活、流量扣减、退款和停复机。
- 不修改 API、DTO、路由或前端。
- 不批量修复历史数据。
- 按用户要求,不新增、修改或运行自动化测试。
## Decisions
### 决策 1数据库先筛选真实孤儿队首再执行 LIMIT
使用 GORM `Raw` 执行 PostgreSQL CTE/窗口查询:
1. 从有效 `status=0` 主套餐按卡/设备载体分组。
2. 每组按 `priority ASC, created_at ASC, id ASC``ROW_NUMBER()=1`
3. 使用相关 `NOT EXISTS` 排除同载体 `status IN (1,2)` 主套餐。
4. 稳定排序后 `LIMIT 100`
删除现有逐条 `Count` 和 Go map 分组。查询仍返回完整 `PackageUsage`,沿用现有投递循环。
**拒绝:扩大 LIMIT。** 只会推迟复现并增加 N+1。
**拒绝:分页遍历所有 pending。** 需要游标状态,复杂度高于一次正确查询。
### 决策 2过期事务提交后再投递现有 Asynq
`processExpiredPackage` 事务只更新旧主套餐和关联加油包。提交成功后,使用普通数据库句柄调用现有 `activateNextPackage` 查询队首并入队。
提交后入队失败时返回错误并记录上下文;同一轮末尾及后续轮询的真实孤儿扫描会再次发现该载体,提供持久状态驱动的补偿。
**拒绝:任务增加固定延迟。** 固定延迟不能证明事务已提交。
**拒绝:引入 Outbox。** 当前线上没有该基础设施,热修不扩张架构。
### 决策 3锁冲突必须触发 Asynq 重试
`ActivateSpecificPackage` 未取得 `RedisPackageActivationLockKey` 时返回现有 `CodePackageActivationConflict`。Handler 原样返回错误,由任务已有 `MaxRetry(3)` 处理。
Redis Key 继续使用 `pkg/constants/redis.go` 的生成函数,不新增硬编码 Key。
### 决策 4显式返回本次是否激活
`ActivateSpecificPackage` 返回 `(bool, error)`
- `true,nil`:本次把待生效套餐推进为生效中;
- `false,nil`:记录已非待生效、存在占位套餐或条件暂不满足;
- `false,error`数据库、Redis 或锁冲突,应由任务重试或记录失败。
Handler 仅在 `true,nil` 时记录“套餐激活成功”。Handler 已在调用前识别 `status=1` 的重复任务并记录幂等跳过。
### 决策 5依赖注入和事务边界保持不变
`PackageActivationHandler` 继续通过结构体字段持有 `*gorm.DB``*redis.Client``*asynq.Client``*ActivationService` 和 Zap Logger不新增单实现接口或工厂。套餐激活仍由 Service 自己开启 GORM 事务Handler 不直接更新新套餐状态。
### 决策 6公共能力与验证
- Audit EventN/A系统自动生命周期推进。
- Domain LedgerN/A`tb_package_usage` 是权威事实。
- Integration LogN/A无新增外部调用。
- OutboxN/A保持当前 Asynq + 周期自愈。
按用户要求不写或运行自动化测试。验证使用 `gofmt``git diff --check``go build ./...`、只读 SQL、查询计划和日志检查。
## Risks / Trade-offs
- **[风险] CTE 扫描大量 pending** → 用 `EXPLAIN (ANALYZE, BUFFERS)` 验证;无证据不新增索引。
- **[风险] 提交后、入队前进程退出** → 下一轮真实孤儿扫描恢复,最长增加一个轮询周期。
- **[风险] 多实例重复入队** → 载体 Redis 锁和套餐状态幂等保证只实际激活一次。
- **[风险] 方法签名变化遗漏调用点** → 使用 `rg` 检查全部调用方并以全量构建证明编译契约。
- **[权衡] Asynq 最终失败后仍依赖轮询重新入队** → 这是当前线上架构的既有补偿边界,本热修不扩建基础设施。
- **[权衡] 不新增自动化测试** → 遵循用户边界以构建、SQL 和日志证据替代。
## Migration Plan
1. 实施真实孤儿查询、提交后入队、锁冲突重试和准确日志。
2. 执行格式化、静态检查、`go build ./...` 和只读 SQL语义检查。
3. 形成独立中文 Lore 热修提交。
4. 部署 Worker观察至少两个轮询周期内真实孤儿收敛和任务日志。
### 回滚
- 无数据库迁移revert 热修提交并重新部署 Worker。
- 已正确激活的套餐保持业务事实,不执行反向 SQL。
- 回滚后新增孤儿继续使用带状态保护的单卡 SQL逐条恢复。
## Open Questions
无。

View File

@@ -0,0 +1,36 @@
## Why
功能 ID`hotfix-main-package-activation-recovery`
线上已第二次出现原主套餐成功过期、队首待生效套餐仍长期停留在 `status=0` 的故障。只读诊断确认当前孤儿扫描固定取出的 100 条记录全部仍有占位套餐,真实孤儿永远无法进入恢复窗口;同时,过期事务提交前投递 Asynq 会让消费者读到旧套餐仍为生效中并把跳过误判为成功。
## What Changes
- PostgreSQL 在 `LIMIT 100` 前完成每个载体队首选择和 `status IN (1,2)` 占位排除,删除逐条检查的 N+1 查询。
- 旧主套餐过期和加油包失效事务提交成功后,才查询队首套餐并投递现有 `package:queue:activation` Asynq 任务。
- Redis 激活锁冲突返回现有套餐激活冲突错误,让 Asynq 按既有策略重试,不再确认假成功。
- Asynq Handler 仅在实际推进套餐或确认已完成幂等事实时记录成功;占位阻塞、等待实名和锁冲突记录明确原因。
- 不新增 Outbox、数据库迁移、依赖或新任务类型保持线上现有纯 Asynq 架构。
- 按用户明确要求,不新增、修改或运行自动化测试;使用全量构建、只读 SQL、查询计划和日志核验。
## Capabilities
### New Capabilities
无。
### Modified Capabilities
- `package-queue-activation`:补充纯 Asynq 过期接续的提交后投递、真实孤儿公平扫描、锁冲突重试和准确成功语义。
## Impact
- **适用范围**:仅当前 `main` 线上代码。
- **架构通道**:旧套餐轮询复杂写用例,沿用 `Polling Handler → GORM transaction → Asynq → Package Activation Service`;不迁移未触碰的套餐模块。
- **代码**`internal/polling/package_activation_handler.go``internal/service/package/activation_service.go`
- **数据库**:无迁移;仅调整 `tb_package_usage` 查询顺序与过滤。
- **API/前端**:无改动。
- **依赖**:继续使用 GORM、PostgreSQL、Redis、Asynq 和 Zap不新增依赖。
- **性能**:孤儿扫描由最多 101 次查询收敛为一次候选查询和有限任务投递;使用查询计划确认数据库耗时。
- **审计与可靠性**Audit Event、Domain Ledger、Integration Log、Outbox 均为 N/A`tb_package_usage` 是权威状态,现有 Asynq + 周期孤儿扫描负责最终恢复。
- **验证**:不写自动化测试;执行 `gofmt`、静态检查、`go build ./...`、只读 SQL及日志核验。

View File

@@ -0,0 +1,77 @@
## MODIFIED Requirements
### Requirement: 当前主套餐过期后自动激活下一个
系统 SHALL 在生效中或已用完主套餐到期时,先提交旧主套餐及关联加油包的状态事务,再投递同一载体队首待生效套餐的现有 Asynq 激活任务;系统 MUST NOT 在旧套餐事务提交前投递任务。
#### Scenario: 事务提交后投递队首套餐
- **WHEN** 轮询处理一个到期的 `status=1``status=2` 主套餐,且存在队首待生效套餐
- **THEN** 系统先提交旧主套餐 `status=3` 和关联加油包 `status=4` 的事务
- **AND** 提交成功后查询稳定队首并投递 `package:queue:activation`
- **AND** 消费者读取时旧主套餐不再处于 `status IN (1,2)`
#### Scenario: 过期事务失败
- **WHEN** 更新旧主套餐或关联加油包失败
- **THEN** 事务回滚
- **AND** 系统不得投递下一套餐激活任务
- **AND** 下一轮过期扫描仍可重新处理
#### Scenario: 提交后入队失败
- **WHEN** 旧套餐事务已经提交,但 Asynq 入队失败或进程退出
- **THEN** 队首套餐保持 `status=0`
- **AND** 周期孤儿扫描重新发现并补投该套餐
#### Scenario: 无待生效套餐
- **WHEN** 旧套餐过期提交后不存在有效队首待生效套餐
- **THEN** 系统不投递激活任务
- **AND** 沿用既有无套餐停机检查
## ADDED Requirements
### Requirement: 孤儿待生效套餐必须公平恢复
系统 SHALL 在数据库中先选择每个卡或设备载体唯一的队首待生效主套餐,并排除仍有 `status IN (1,2)` 占位主套餐的载体,最后才执行单轮 100 个真实孤儿上限。队首顺序 MUST 为 `priority ASC, created_at ASC, id ASC`
#### Scenario: 固定窗口全部为非孤儿
- **WHEN** 排序靠前的 100 条待生效记录均有占位主套餐,窗口之后存在真实孤儿
- **THEN** 数据库先排除前 100 条非孤儿
- **AND** 窗口后的真实孤儿进入本轮候选并被投递
#### Scenario: 同一载体有多条 pending
- **WHEN** 一个真实孤儿载体存在多条待生效主套餐
- **THEN** 本轮只选择 priority 最小、created_at 最早、id 最小的一条
- **AND** 该载体只占一个恢复名额
#### Scenario: 生效中或已用完套餐占位
- **WHEN** pending 所属载体存在 `status=1``status=2` 主套餐
- **THEN** 该载体不得进入孤儿候选
### Requirement: Asynq 激活结果必须准确且可重试
系统 SHALL 仅在本次实际把套餐推进为 `status=1` 时记录新激活成功。Redis 激活锁冲突 MUST 返回错误,使 Asynq 按既有 `MaxRetry(3)` 重试,不得确认未执行任务成功。
#### Scenario: 锁冲突触发重试
- **WHEN** 消费者未取得载体级套餐激活锁
- **THEN** Service 返回套餐激活冲突错误
- **AND** Handler 将错误返回 Asynq
- **AND** 本次不记录激活成功
#### Scenario: 本次实际激活
- **WHEN** 套餐为待生效、载体无占位主套餐且满足激活条件
- **THEN** Service 返回 `activated=true`
- **AND** Handler 记录包含套餐使用记录和触发类型的成功日志
#### Scenario: 条件暂不满足
- **WHEN** Service 复检发现占位套餐或等待实名条件
- **THEN** Service 返回 `activated=false`且不修改套餐
- **AND** Handler 记录未激活原因,不记录成功

View File

@@ -0,0 +1,20 @@
## 1. Main 分支隔离与基线
- [x] 1.1 确认变更仅面向 `main` 的纯 Asynq 套餐接续链路,并保持 Outbox 与新任务基础设施为非目标。
- [x] 1.2 检索套餐激活方法的全部调用点及现有错误码、Redis 键和任务重试配置,锁定最小修改边界。
## 2. 过期接续与孤儿恢复
- [x] 2.1 将孤儿恢复改为数据库先按载体选择队首并排除 `status IN (1,2)` 占位套餐,最后限制 100 个真实孤儿,删除逐条占位查询。
- [x] 2.2 将旧主套餐过期后的下一套餐投递移到事务提交后,保留现有停机检查和周期孤儿补偿。
## 3. 激活结果与重试语义
- [x] 3.1 让指定套餐激活显式返回是否实际激活,并在 Redis 激活锁冲突时返回现有套餐激活冲突错误。
- [x] 3.2 调整 Asynq Handler 日志,仅在实际激活时记录成功,未激活时记录明确的跳过信息。
## 4. 文档、验证与提交
- [x] 4.1 更新功能总结和 README说明 `main` 纯 Asynq 热修边界、部署观察项及回滚方式。
- [x] 4.2 执行 `gofmt``git diff --check``go build ./...` 和 OpenSpec 严格校验;按用户要求不新增、修改或运行自动化测试。
- [x] 4.3 仅暂存本热修代码、文档和独立 OpenSpec创建符合 Lore 协议的中文提交。

View File

@@ -97,13 +97,38 @@
- **WHEN** 当前账号请求资金概况列表但未提供 `shop_id`
- **THEN** 系统继续按既有分页、数据范围、店铺名称和主账号用户名条件返回结果
### Requirement: 历史待审批线下代理充值可主动接入企业微信审批
系统 SHALL 提供 `POST /api/admin/agent-recharges/{id}/trigger-approval`,使具有既有代理充值管理访问权限的后台账号可为历史线下代理充值申请主动创建企业微信审批。系统 MUST 仅在线下充值记录处于待审批状态且 `approval_instance_id` 为空时创建审批;审批发起人 MUST 使用该充值记录的原创建账号。创建成功后,系统 MUST 原子保存唯一审批实例、审批提交请求及充值记录的审批实例关联,并返回更新后的充值申请审批摘要。
#### Scenario: 主动发起历史线下代理充值审批成功
- **GIVEN** 线下代理充值申请处于待审批状态、未关联审批实例,且其原创建账号和企业微信线下充值审批场景均可用
- **WHEN** 有既有代理充值管理访问权限的后台账号请求 `POST /api/admin/agent-recharges/{id}/trigger-approval`
- **THEN** 系统创建以原创建账号为发起人的唯一企业微信审批并返回审批摘要,后续由既有可靠提交流程提交至企业微信
#### Scenario: 在线、非待审批或已发起记录被拒绝
- **WHEN** 请求主动发起的充值记录不是线下充值、不是待审批状态或已关联审批实例
- **THEN** 系统返回状态冲突且不创建新的审批实例或提交请求
#### Scenario: 并发主动发起同一充值审批
- **WHEN** 两个请求同时为同一符合条件的线下代理充值申请主动发起审批
- **THEN** 系统至多创建一个审批实例和一个审批提交请求,未成功创建关联的请求返回冲突
#### Scenario: 原创建人或审批渠道不可用
- **WHEN** 充值申请原创建账号不可用,或企业微信线下充值审批场景不可用
- **THEN** 系统返回相应错误,充值申请保持未关联审批实例,修复条件后可再次发起
## 可达操作索引
本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。
### 代理预充值
`GET /api/admin/agent-recharges`(查询代理充值订单列表);`POST /api/admin/agent-recharges`(创建代理充值订单);`GET /api/admin/agent-recharges/{id}`(查询代理充值订单详情);`POST /api/admin/agent-recharges/{id}/offline-pay`(确认线下充值);`GET /api/admin/agent-recharges/{id}/payment-status`(查询代理充值本地支付与到账状态);`POST /api/admin/agent-recharges/{id}/reject`(驳回代理充值订单);`GET /api/admin/agent-recharges/payment-methods`(查询代理在线充值可用支付方式)。
`GET /api/admin/agent-recharges`(查询代理充值订单列表);`POST /api/admin/agent-recharges`(创建代理充值订单);`GET /api/admin/agent-recharges/{id}`(查询代理充值订单详情);`POST /api/admin/agent-recharges/{id}/trigger-approval`(补发历史线下代理充值审批);`POST /api/admin/agent-recharges/{id}/offline-pay`(确认线下充值);`GET /api/admin/agent-recharges/{id}/payment-status`(查询代理充值本地支付与到账状态);`POST /api/admin/agent-recharges/{id}/reject`(驳回代理充值订单);`GET /api/admin/agent-recharges/payment-methods`(查询代理在线充值可用支付方式)。
### 代理商资金管理

View File

@@ -38,6 +38,31 @@
- **WHEN** 该代理查询退款列表或退款申请详情
- **THEN** 系统返回空列表或不存在,且不泄露任何退款申请
### Requirement: 历史待审批退款可主动接入企业微信审批
系统 SHALL 提供 `POST /api/admin/refunds/{id}/trigger-approval`,使具有既有退款管理访问权限的后台账号可为历史退款申请主动创建企业微信审批。系统 MUST 仅在退款申请状态为待审批且 `approval_instance_id` 为空时创建审批;审批发起人 MUST 使用该退款申请的原创建账号。创建成功后,系统 MUST 原子保存唯一审批实例、审批提交请求及退款申请的审批实例关联,并返回更新后的退款申请审批摘要。
#### Scenario: 主动发起历史退款审批成功
- **GIVEN** 退款申请处于待审批状态、未关联审批实例,且其原创建账号和企业微信退款审批场景均可用
- **WHEN** 有既有退款管理访问权限的后台账号请求 `POST /api/admin/refunds/{id}/trigger-approval`
- **THEN** 系统创建以原创建账号为发起人的唯一企业微信审批并返回审批摘要,后续由既有可靠提交流程提交至企业微信
#### Scenario: 非待审批或已发起记录被拒绝
- **WHEN** 请求主动发起的退款申请不是待审批状态或已关联审批实例
- **THEN** 系统返回状态冲突且不创建新的审批实例或提交请求
#### Scenario: 并发主动发起同一退款审批
- **WHEN** 两个请求同时为同一符合条件的退款申请主动发起审批
- **THEN** 系统至多创建一个审批实例和一个审批提交请求,未成功创建关联的请求返回冲突
#### Scenario: 原创建人或审批渠道不可用
- **WHEN** 退款申请原创建账号不可用,或企业微信退款审批场景不可用
- **THEN** 系统返回相应错误,退款申请保持未关联审批实例,修复条件后可再次发起
## 可达操作索引
本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。
@@ -48,7 +73,7 @@
### 退款管理
`GET /api/admin/refunds`(退款申请列表);`POST /api/admin/refunds`(创建退款申请);`GET /api/admin/refunds/{id}`(退款申请详情);`POST /api/admin/refunds/{id}/approve`(审批通过退款申请);`POST /api/admin/refunds/{id}/reject`(审批拒绝退款申请);`POST /api/admin/refunds/{id}/resubmit`(重新提交退款申请);`POST /api/admin/refunds/{id}/return`(退回退款申请)。
`GET /api/admin/refunds`(退款申请列表);`POST /api/admin/refunds`(创建退款申请);`GET /api/admin/refunds/{id}`(退款申请详情);`POST /api/admin/refunds/{id}/trigger-approval`(补发历史退款审批);`POST /api/admin/refunds/{id}/approve`(审批通过退款申请);`POST /api/admin/refunds/{id}/reject`(审批拒绝退款申请);`POST /api/admin/refunds/{id}/resubmit`(重新提交退款申请);`POST /api/admin/refunds/{id}/return`(退回退款申请)。
### 换货管理