13 Commits

Author SHA1 Message Date
395e5fb47c 修复套餐接续停机竞态与轮询兜底
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 11m13s
2026-08-29 16:27:48 +08:00
62f3d25e81 修复批量订购 2026-08-26 14:54:26 +08:00
5797fd0e94 修复企微兜底机制
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 10m40s
2026-08-20 18:06:45 +08:00
ba677a35e1 Update deploy.yaml
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Failing after 28h53m35s
2026-08-20 12:06:58 +08:00
22b95db2f9 每日清理 2026-08-20 12:03:17 +08:00
143df60485 修复 2026-08-18 17:32:15 +08:00
656a921ff0 富友支付支持 2026-08-18 17:13:20 +08:00
46c8e819df 修改相关证据 2026-08-18 16:30:16 +08:00
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
101 changed files with 8907 additions and 16912 deletions

View File

@@ -3,9 +3,7 @@ name: 构建并部署到测试环境(无 SSH
on:
push:
branches:
- Iteration/7-11
- dev
- test
- main
env:
REGISTRY: registry.boss160.cn
@@ -30,15 +28,7 @@ jobs:
- name: 设置镜像标签
id: tag
run: |
if [ "${{ github.ref }}" = "refs/heads/Iteration/7-11" ]; then
echo "tag=latest" >> $GITHUB_OUTPUT
elif [ "${{ github.ref }}" = "refs/heads/dev" ]; then
echo "tag=dev" >> $GITHUB_OUTPUT
elif [ "${{ github.ref }}" = "refs/heads/test" ]; then
echo "tag=test" >> $GITHUB_OUTPUT
else
echo "tag=unknown" >> $GITHUB_OUTPUT
fi
- name: 登录 Docker Registry
run: |
@@ -61,8 +51,8 @@ jobs:
docker push ${{ env.WORKER_IMAGE }}:${{ steps.tag.outputs.tag }}
docker push ${{ env.WORKER_IMAGE }}:${{ github.sha }}
- name: 部署到本地(仅 Iteration/7-11 分支)
if: github.ref == 'refs/heads/Iteration/7-11'
- name: 部署到测试环境(仅 main 分支)
if: github.ref == 'refs/heads/main'
run: |
# 确保部署目录存在(仅需日志目录,配置已嵌入二进制文件)
mkdir -p ${{ env.DEPLOY_DIR }}/logs

2
.gitignore vendored
View File

@@ -111,6 +111,4 @@ scripts/batch_package_purchase/assets.example_购买结果_20260715_115804.csv
scripts/batch_package_purchase/assets.example_购买结果_20260715_115814.csv
scripts/migration/output
# LongHorizon 本地执行证据
.lh-harness/
.scratch/go-build-cache

File diff suppressed because it is too large Load Diff

View File

@@ -41,7 +41,7 @@
- 不访问生产服务、真实支付渠道或外部审批系统进行自动验证。
- 生产环境为 systemd 管理的手工二进制发布,和仓库 Docker/CI 测试环境不同;生产发布、迁移与回滚事实见 [`docs/deployment/production-runbook.md`](docs/deployment/production-runbook.md)。Agent 不连接生产主机或数据库,生产操作由维护者执行并提供结果。
- 不把密钥、Token、证书或个人敏感数据写入代码、文档和日志。
- `.lh-harness/` 仅保存本地执行证据,不是事实源且不得纳入 Git
- `docs/verification/context-reset/` 仅保存上下文健康检查证据,不是业务事实源;证据应随项目文档维护
- 自动化测试当前为 N/A用户决策不恢复旧测试也不写虚假测试入口。
## 架构选择

1041
README.md

File diff suppressed because it is too large Load Diff

View File

@@ -1,4 +1,4 @@
// Command audit-retention-simulate 在测试环境演练完整自然月归档与清理边界。
// Command audit-retention-simulate 在测试环境演练逐日归档与清理边界。
package main
import (
@@ -6,22 +6,21 @@ import (
"crypto/sha256"
"fmt"
"os"
"path/filepath"
"strings"
"time"
"github.com/bytedance/sonic"
"github.com/hibiken/asynq"
"go.uber.org/zap"
"gorm.io/datatypes"
"gorm.io/gorm"
auditarchive "github.com/break/junhong_cmp_fiber/internal/application/auditarchive"
auditinfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
"github.com/break/junhong_cmp_fiber/internal/model"
taskapp "github.com/break/junhong_cmp_fiber/internal/task"
"github.com/break/junhong_cmp_fiber/pkg/config"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/database"
logpkg "github.com/break/junhong_cmp_fiber/pkg/logger"
"github.com/break/junhong_cmp_fiber/pkg/storage"
)
@@ -67,6 +66,11 @@ func run(ctx context.Context) error {
return fmt.Errorf("仅允许显式确认的测试数据库,当前数据库为 %q", cfg.Database.DBName)
}
logger := zap.NewNop()
if err := os.MkdirAll(filepath.Dir(cfg.Logging.RetentionLog.Filename), 0755); err != nil {
return err
}
retentionLogger, syncRetentionLogger := logpkg.NewRetentionLogger(cfg.Logging.Level, logpkg.LogRotationConfig{Filename: cfg.Logging.RetentionLog.Filename, MaxSize: cfg.Logging.RetentionLog.MaxSize, MaxBackups: cfg.Logging.RetentionLog.MaxBackups, MaxAge: cfg.Logging.RetentionLog.MaxAge, Compress: cfg.Logging.RetentionLog.Compress})
defer func() { _ = syncRetentionLogger() }()
db, err := database.InitPostgreSQL(&cfg.Database, logger)
if err != nil {
return err
@@ -81,8 +85,7 @@ func run(ctx context.Context) error {
if err != nil {
return err
}
auditWriter := auditinfra.NewWriter(auditinfra.NewRegistry(), nil)
service, err := auditarchive.NewService(db, provider, simulationInstance, auditWriter)
service, err := auditarchive.NewService(db, provider, simulationInstance)
if err != nil {
return err
}
@@ -108,14 +111,17 @@ func run(ctx context.Context) error {
}
}
payload, _ := sonic.Marshal(taskapp.AuditMonthlyRetentionPayload{ArchiveMonth: simulationMonth})
dryRunTask := asynq.NewTask(constants.TaskTypeAuditMonthlyRetention, payload)
if err := taskapp.NewAuditMonthlyRetentionHandler(service, logger, false).Handle(ctx, dryRunTask); err != nil {
return fmt.Errorf("月度只读演练失败: %w", err)
var dryRun auditarchive.RetentionResult
for date := monthStart; date.Before(monthEnd); date = date.AddDate(0, 0, 1) {
result, retainErr := service.RetainDate(ctx, date, false)
if retainErr != nil {
return fmt.Errorf("逐日只读演练失败: %w", retainErr)
}
dryRun, err := service.ValidateMonth(ctx, monthStart)
if err != nil {
return err
retentionLogger.Info("日留存仿真只读校验通过", zap.String("archive_date", result.ArchiveDate))
dryRun.EventCount += result.EventCount
dryRun.ResourceCount += result.ResourceCount
dryRun.IntegrationCount += result.IntegrationCount
dryRun.EstimatedBatches += result.EstimatedBatches
}
summary, err := collectBeforeCleanup(ctx, db, monthStart, monthEnd, dryRun)
if err != nil {
@@ -125,16 +131,39 @@ func run(ctx context.Context) error {
return fmt.Errorf("只归档模式写入了 %d 个清理断点", summary.CleanupMarkersBefore)
}
cleanupTask := asynq.NewTask(constants.TaskTypeAuditMonthlyRetention, payload)
if err := taskapp.NewAuditMonthlyRetentionHandler(service, logger, true).Handle(ctx, cleanupTask); err != nil {
return fmt.Errorf("隔离测试库物理清理演练失败: %w", err)
for date := monthStart; date.Before(monthEnd); date = date.AddDate(0, 0, 1) {
result, retainErr := service.RetainDate(ctx, date, true)
if retainErr != nil {
return fmt.Errorf("隔离测试库逐日物理清理演练失败: %w", retainErr)
}
retentionLogger.Info("日留存仿真物理清理完成", zap.String("archive_date", result.ArchiveDate))
}
if err := collectAfterCleanup(ctx, db, monthStart, monthEnd, &summary); err != nil {
return err
}
if summary.TargetRowsAfterCleanup != 0 || summary.BoundaryRowsAfterCleanup != 6 || summary.CleanupMarkersAfter != int64(summary.Days*2) || !summary.RetentionAuditRecorded {
return fmt.Errorf("清理范围复核失败: target=%d boundary=%d markers=%d audit=%t",
summary.TargetRowsAfterCleanup, summary.BoundaryRowsAfterCleanup, summary.CleanupMarkersAfter, summary.RetentionAuditRecorded)
if summary.TargetRowsAfterCleanup != 1 || summary.BoundaryRowsAfterCleanup != 6 || summary.CleanupMarkersAfter != int64(summary.Days*2) {
return fmt.Errorf("清理范围复核失败: target=%d boundary=%d markers=%d", summary.TargetRowsAfterCleanup, summary.BoundaryRowsAfterCleanup, summary.CleanupMarkersAfter)
}
if err := db.WithContext(ctx).Model(&model.IntegrationLog{}).
Where("integration_id = ?", "int_"+simulationPrefix+"-"+monthStart.Format("20060102")).
Updates(map[string]any{"result": constants.IntegrationResultSuccess, "updated_at": time.Now()}).Error; err != nil {
return err
}
if _, err := service.RetainDate(ctx, monthStart, true); err != nil {
return fmt.Errorf("pending 终结后的续跑清理演练失败: %w", err)
}
var remaining int64
if err := db.WithContext(ctx).Model(&model.IntegrationLog{}).Where("created_at >= ? AND created_at < ?", monthStart, monthEnd).Count(&remaining).Error; err != nil {
return err
}
if remaining != 0 {
return fmt.Errorf("pending 终结后的续跑清理未完成,剩余 %d 条记录", remaining)
}
if err := syncRetentionLogger(); err != nil {
return fmt.Errorf("刷新日留存仿真日志失败: %w", err)
}
if _, err := os.Stat(cfg.Logging.RetentionLog.Filename); err != nil {
return fmt.Errorf("独立日留存日志未生成: %w", err)
}
encoded, _ := sonic.MarshalIndent(summary, "", " ")
fmt.Println(string(encoded))
@@ -256,12 +285,7 @@ func archiveMonth(ctx context.Context, db *gorm.DB, service *auditarchive.Servic
if err := service.ArchiveDate(ctx, start); err != nil {
return fmt.Errorf("Audit 故障重试演练失败: %w", err)
}
if err := db.WithContext(ctx).Model(&model.IntegrationLog{}).
Where("integration_id = ?", "int_"+simulationPrefix+"-"+start.Format("20060102")).
Updates(map[string]any{"result": constants.IntegrationResultSuccess, "updated_at": time.Now()}).Error; err != nil {
return err
}
return service.FinalizeIntegrationMonth(ctx, start)
return nil
}
func collectBeforeCleanup(ctx context.Context, db *gorm.DB, start, end time.Time, dryRun auditarchive.RetentionResult) (simulationSummary, error) {
@@ -311,10 +335,5 @@ func collectAfterCleanup(ctx context.Context, db *gorm.DB, start, end time.Time,
Count(&summary.CleanupMarkersAfter).Error; err != nil {
return err
}
var auditCount int64
if err := db.WithContext(ctx).Model(&model.AuditEvent{}).Where("event_id = ?", "evt_retention_"+strings.ReplaceAll(simulationMonth, "-", "_")).Count(&auditCount).Error; err != nil {
return err
}
summary.RetentionAuditRecorded = auditCount == 1
return nil
}

View File

@@ -122,6 +122,11 @@ func runWorker(cfg *config.Config) {
defer func() {
_ = logger.Sync() // 忽略 sync 错误
}()
retentionLogger, syncRetentionLogger := logger.NewRetentionLogger(cfg.Logging.Level, logger.LogRotationConfig{
Filename: cfg.Logging.RetentionLog.Filename, MaxSize: cfg.Logging.RetentionLog.MaxSize,
MaxBackups: cfg.Logging.RetentionLog.MaxBackups, MaxAge: cfg.Logging.RetentionLog.MaxAge, Compress: cfg.Logging.RetentionLog.Compress,
})
defer func() { _ = syncRetentionLogger() }()
appLogger := logger.GetAppLogger()
ctx, cancel := context.WithCancel(context.Background())
@@ -148,7 +153,7 @@ func runWorker(cfg *config.Config) {
taskHandler.RegisterHandlers()
registerWeComApprovalTasks(taskHandler.GetMux(), runtime, cfg, appLogger)
registerAgentRechargeRecoveryTask(taskHandler.GetMux(), runtime, appLogger)
registerAuditArchiveTask(taskHandler.GetMux(), runtime, cfg.Worker.AuditRetentionCleanupEnabled, appLogger)
registerAuditArchiveTask(taskHandler.GetMux(), runtime, cfg.Worker.AuditRetentionCleanupEnabled, appLogger, retentionLogger)
outboxHandler := outbox.NewHandler(runtime.outboxConsumers)
taskHandler.GetMux().HandleFunc(constants.TaskTypeOutboxDeliver, outboxHandler.Handle)
startOutboxRelay(ctx, runtime, cfg.Worker.InstanceName, appLogger)
@@ -444,6 +449,7 @@ func registerAgentRechargeRecoveryTask(mux *asynq.ServeMux, runtime *workerRunti
runtime.db,
paymentInfra.NewWechatWebAdapter(wechat.NewRedisCache(runtime.redisClient), integration, appLogger),
paymentInfra.NewAlipayWapAdapter(integration, appLogger),
paymentInfra.NewFuiouScanAdapter(integration, appLogger),
confirm,
runtime.workerResult.Services.PaymentAudit,
)
@@ -793,43 +799,32 @@ func registerAsynqScheduleTasks(asynqScheduler *asynq.Scheduler) error {
)); err != nil {
return fmt.Errorf("注册 Integration Log 每日冷归档定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 5 1 * *", asynq.NewTask(
constants.TaskTypeIntegrationMonthlyFinalize,
nil,
asynq.MaxRetry(10),
asynq.Timeout(6*time.Hour),
asynq.Unique(27*24*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeIntegrationMonthlyFinalize)),
)); err != nil {
return fmt.Errorf("注册 Integration Log 月度最终版本复核任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 6 1 * *", asynq.NewTask(
constants.TaskTypeAuditMonthlyRetention,
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 5 * * *", asynq.NewTask(
constants.TaskTypeAuditDailyRetention,
nil,
asynq.MaxRetry(10),
asynq.Timeout(12*time.Hour),
asynq.Unique(27*24*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAuditMonthlyRetention)),
asynq.Unique(23*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAuditDailyRetention)),
)); err != nil {
return fmt.Errorf("注册月度日志留存演练或清理任务失败: %w", err)
return fmt.Errorf("注册日志留存演练或清理任务失败: %w", err)
}
return nil
}
// registerAuditArchiveTask 注册 Audit 与 Integration 冷归档任务处理器。
func registerAuditArchiveTask(mux *asynq.ServeMux, runtime *workerRuntime, cleanupEnabled bool, appLogger *zap.Logger) {
func registerAuditArchiveTask(mux *asynq.ServeMux, runtime *workerRuntime, cleanupEnabled bool, appLogger, retentionLogger *zap.Logger) {
if runtime.storageSvc == nil {
appLogger.Warn("对象存储未配置,审计归档任务将在执行时重试")
mux.HandleFunc(constants.TaskTypeAuditDailyArchive, task.NewAuditDailyArchiveHandler(nil, appLogger).Handle)
integrationHandler := task.NewIntegrationArchiveHandler(nil, appLogger)
mux.HandleFunc(constants.TaskTypeIntegrationDailyArchive, integrationHandler.HandleDaily)
mux.HandleFunc(constants.TaskTypeIntegrationMonthlyFinalize, integrationHandler.HandleMonthlyFinalize)
mux.HandleFunc(constants.TaskTypeAuditMonthlyRetention, task.NewAuditMonthlyRetentionHandler(nil, appLogger, cleanupEnabled).Handle)
mux.HandleFunc(constants.TaskTypeAuditDailyRetention, task.NewAuditRetentionHandler(nil, retentionLogger, cleanupEnabled).Handle)
return
}
auditWriter, ok := runtime.workerResult.Services.PaymentAudit.(*auditInfra.Writer)
if !ok || auditWriter == nil {
appLogger.Fatal("初始化月度日志留存清理失败:统一审计 Writer 未配置")
appLogger.Fatal("初始化日志留存清理失败:统一审计 Writer 未配置")
}
service, err := auditArchiveApp.NewService(runtime.db, runtime.storageSvc.Provider(), constants.AuditArchiveInstanceID, auditWriter)
if err != nil {
@@ -838,13 +833,11 @@ func registerAuditArchiveTask(mux *asynq.ServeMux, runtime *workerRuntime, clean
mux.HandleFunc(constants.TaskTypeAuditDailyArchive, task.NewAuditDailyArchiveHandler(service, appLogger).Handle)
integrationHandler := task.NewIntegrationArchiveHandler(service, appLogger)
mux.HandleFunc(constants.TaskTypeIntegrationDailyArchive, integrationHandler.HandleDaily)
mux.HandleFunc(constants.TaskTypeIntegrationMonthlyFinalize, integrationHandler.HandleMonthlyFinalize)
mux.HandleFunc(constants.TaskTypeAuditMonthlyRetention, task.NewAuditMonthlyRetentionHandler(service, appLogger, cleanupEnabled).Handle)
mux.HandleFunc(constants.TaskTypeAuditDailyRetention, task.NewAuditRetentionHandler(service, retentionLogger, cleanupEnabled).Handle)
appLogger.Info("注册审计归档任务处理器",
zap.String("audit_task_type", constants.TaskTypeAuditDailyArchive),
zap.String("integration_daily_task_type", constants.TaskTypeIntegrationDailyArchive),
zap.String("integration_monthly_task_type", constants.TaskTypeIntegrationMonthlyFinalize),
zap.String("retention_task_type", constants.TaskTypeAuditMonthlyRetention),
zap.String("retention_task_type", constants.TaskTypeAuditDailyRetention),
zap.Bool("retention_cleanup_enabled", cleanupEnabled))
}

View File

@@ -74,6 +74,13 @@ DB_PASSWORD='<密码>' DB_NAME=<库名> DB_SSLMODE=<模式> \
迁移失败时不启动新二进制;按失败迁移的事务状态决定处理,必要时恢复已确认可用的数据库备份。启动失败时覆盖回部署前备份的二进制,再恢复数据库备份(如迁移已改变数据库)。
### 零金额退款发布后核验
发布本次退款审批变更后,维护者应先等待既有重试处理稳定事件 `approval:26:approved`;若重试已耗尽,按受控运维流程重放同一事件,不得直接修改退款、订单或钱包数据。随后核验:
1. 退款单 `RF20260820170954264700` 已通过,关联订单支付状态为已退款。
2. 该退款没有代理主钱包或资产钱包的零金额回款流水。
### 锁的含义与发布影响
`000171` 会对 `tb_agent_wallet``000178` 会对 `tb_iot_card` 使用 PostgreSQL `ACCESS EXCLUSIVE` 锁。该锁执行期间会阻塞该表的读写及其他 DDL直到迁移事务提交或回滚若有未结束业务查询/事务,它也会等待。因此必须在 API 和全部 Worker 停止后执行,并在迁移前检查没有长事务。锁持续时间取决于表数据量、索引创建和等待中的旧事务;不能从仓库估算具体秒数。
@@ -88,3 +95,12 @@ DB_PASSWORD='<密码>' DB_NAME=<库名> DB_SSLMODE=<模式> \
## 已确认数据库备份
每日凌晨 02:00 自动备份 `junhong_cmp_prod`:数据库运行在 Docker 容器 `postgres` 中,备份脚本执行 `pg_dump -Fc -Z 6`,写入 `/data/backups/postgresql/<库名>_<时间>.dump`,同时生成 MD5 文件并以 `pg_restore --list` 校验结构;保留 30 天。发布前仍须人工新建一次备份并确认校验通过,不能只依赖凌晨的最近备份。恢复命令待维护者实际演练或确认后补充。
## 日审计日志留存启用与恢复
日留存由调度 Worker 每日处理 Asia/Shanghai 的昨天及更早连续积压日期;`JUNHONG_WORKER_AUDIT_RETENTION_CLEANUP_ENABLED=false` 时仅验证归档、对象、manifest、数据库数量和日期连续性不写清理断点、不删除在线数据。
1. 发布新 Worker 后保持开关关闭。维护者先核对 `tb_log_archive_run` 中 Audit 与 Integration 两个来源从历史最早在线日期起没有缺失账本;缺失日期必须先受控补归档,不能跳过失败日。
2. 以关闭开关的 Worker 完成只读演练,观察 `/opt/junhong_cmp/worker/logs/audit-retention.log`:每个日期应有来源计数和耗时;若出现日期、来源、失败分类或安全错误摘要,先修复该日归档或 pending 记录后再演练。
3. 低峰期将调度 Worker 的 `JUNHONG_WORKER_AUDIT_RETENTION_CLEANUP_ENABLED=true`,重启该 Worker并持续观察上述独立日志、`tb_log_archive_run.cleanup_started_at` / `cleaned_at` 断点及表大小。任何日期失败都会阻断该日和后续日期,不得手工跳过。
4. 异常时立即将开关改回 `false` 并重启调度 Worker。已清理日期按已验证的对象和 manifest 执行归档恢复;尚未清理或阻断的日期仍保留在线,无需数据库恢复。恢复后先重新执行只读演练,再决定是否重新开启清理。

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。

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

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

@@ -49,17 +49,18 @@ type OnlineCreationService struct {
db *gorm.DB
wechat OnlinePaymentPort
alipay OnlinePaymentPort
fuiou OnlinePaymentPort
audit PaymentAuditWriter
}
// NewOnlineCreationService 创建代理在线充值用例并以结构体字段注入个渠道 Adapter。
func NewOnlineCreationService(db *gorm.DB, wechat, alipay OnlinePaymentPort, audit PaymentAuditWriter) *OnlineCreationService {
return &OnlineCreationService{db: db, wechat: wechat, alipay: alipay, audit: audit}
// NewOnlineCreationService 创建代理在线充值用例并以结构体字段注入个渠道 Adapter。
func NewOnlineCreationService(db *gorm.DB, wechat, alipay, fuiou OnlinePaymentPort, audit PaymentAuditWriter) *OnlineCreationService {
return &OnlineCreationService{db: db, wechat: wechat, alipay: alipay, fuiou: fuiou, audit: audit}
}
// Execute 以短事务建单,事务外生成支付链接,再条件保存链接或关闭失败订单。
func (s *OnlineCreationService) Execute(ctx context.Context, command CreateOnlineCommand) (*CreateOnlineResult, error) {
if s == nil || s.db == nil || s.wechat == nil || s.alipay == nil || s.audit == nil {
if s == nil || s.db == nil || s.wechat == nil || s.alipay == nil || s.fuiou == nil || s.audit == nil {
return nil, apperrors.New(apperrors.CodeServiceUnavailable, "代理在线充值能力未配置")
}
command.PaymentMethod = strings.TrimSpace(command.PaymentMethod)
@@ -139,7 +140,7 @@ func (s *OnlineCreationService) AvailablePaymentMethods(ctx context.Context, use
}
return result, apperrors.Wrap(apperrors.CodeDatabaseError, err, "查询生效支付配置失败")
}
if s.wechat.Available(&config) {
if s.wechat.Available(&config) || s.fuiou.Available(&config) {
result.Methods = append(result.Methods, constants.RechargeMethodWechat)
}
if s.alipay.Available(&config) {
@@ -180,7 +181,7 @@ func (s *OnlineCreationService) loadCreationFacts(
}
return nil, nil, nil, nil, nil, apperrors.Wrap(apperrors.CodeDatabaseError, err, "查询生效支付配置失败")
}
adapter := s.adapter(command.PaymentMethod)
adapter := s.adapter(command.PaymentMethod, &config)
if adapter == nil || !adapter.Available(&config) {
return nil, nil, nil, nil, nil, apperrors.New(apperrors.CodeNoPaymentConfig)
}
@@ -209,7 +210,7 @@ func (s *OnlineCreationService) createLocalFacts(
expireMinutes = model.DefaultAliPayExpireMinutes
}
expireAt := time.Now().Add(time.Duration(expireMinutes) * time.Minute)
channel, requestID := command.PaymentMethod, command.RequestID
channel, requestID := paymentChannel(command.PaymentMethod, config), command.RequestID
record := &model.AgentRechargeRecord{
UserID: account.ID, AgentWalletID: wallet.ID, ShopID: shop.ID, RechargeNo: rechargeNo,
Amount: command.Amount, PaymentMethod: command.PaymentMethod, PaymentChannel: &channel,
@@ -247,6 +248,9 @@ func paymentMerchantIdentity(paymentMethod string, config *model.WechatConfig) s
return ""
}
if paymentMethod == constants.RechargeMethodWechat {
if config.ProviderType == model.ProviderTypeFuiou {
return config.FyMchntCd
}
return config.WxMchID
}
if paymentMethod == constants.RechargeMethodAlipay {
@@ -255,6 +259,14 @@ func paymentMerchantIdentity(paymentMethod string, config *model.WechatConfig) s
return ""
}
// paymentChannel 返回实际支付渠道:富友配置下微信业务方式落库为 fuiou其余与业务方式一致。
func paymentChannel(paymentMethod string, config *model.WechatConfig) string {
if paymentMethod == constants.RechargeMethodWechat && config != nil && config.ProviderType == model.ProviderTypeFuiou {
return model.ProviderTypeFuiou
}
return paymentMethod
}
func (s *OnlineCreationService) loadReplay(
ctx context.Context,
command CreateOnlineCommand,
@@ -316,8 +328,11 @@ func (s *OnlineCreationService) closeFailedCreation(ctx context.Context, result
})
}
func (s *OnlineCreationService) adapter(paymentMethod string) OnlinePaymentPort {
func (s *OnlineCreationService) adapter(paymentMethod string, config *model.WechatConfig) OnlinePaymentPort {
if paymentMethod == constants.RechargeMethodWechat {
if config != nil && config.ProviderType == model.ProviderTypeFuiou {
return s.fuiou
}
return s.wechat
}
if paymentMethod == constants.RechargeMethodAlipay {

View File

@@ -16,19 +16,20 @@ type RecoverOnlinePaymentService struct {
db *gorm.DB
wechat OnlinePaymentPort
alipay OnlinePaymentPort
fuiou OnlinePaymentPort
confirm *ConfirmOnlinePaymentService
audit PaymentAuditWriter
now func() time.Time
}
// NewRecoverOnlinePaymentService 创建代理在线充值支付恢复用例。
func NewRecoverOnlinePaymentService(db *gorm.DB, wechat, alipay OnlinePaymentPort, confirm *ConfirmOnlinePaymentService, audit PaymentAuditWriter) *RecoverOnlinePaymentService {
return &RecoverOnlinePaymentService{db: db, wechat: wechat, alipay: alipay, confirm: confirm, audit: audit, now: time.Now}
func NewRecoverOnlinePaymentService(db *gorm.DB, wechat, alipay, fuiou OnlinePaymentPort, confirm *ConfirmOnlinePaymentService, audit PaymentAuditWriter) *RecoverOnlinePaymentService {
return &RecoverOnlinePaymentService{db: db, wechat: wechat, alipay: alipay, fuiou: fuiou, confirm: confirm, audit: audit, now: time.Now}
}
// ProcessBatch 按固定批次读取本地待处理事实并调用对应渠道收敛状态。
func (s *RecoverOnlinePaymentService) ProcessBatch(ctx context.Context) (int, error) {
if s == nil || s.db == nil || s.wechat == nil || s.alipay == nil || s.confirm == nil || s.audit == nil {
if s == nil || s.db == nil || s.wechat == nil || s.alipay == nil || s.fuiou == nil || s.confirm == nil || s.audit == nil {
return 0, errors.New(errors.CodeServiceUnavailable, "代理在线充值支付恢复能力未配置")
}
now := s.now().UTC()
@@ -70,7 +71,7 @@ func (s *RecoverOnlinePaymentService) ProcessBatch(ctx context.Context) (int, er
}
func (s *RecoverOnlinePaymentService) recoverOne(ctx context.Context, payment *model.Payment, recharge *model.AgentRechargeRecord, config *model.WechatConfig, now time.Time) error {
adapter := s.adapter(payment.PaymentMethod)
adapter := s.adapter(payment.PaymentMethod, config)
if adapter == nil {
return errors.New(errors.CodeNoPaymentConfig, "代理充值创建时支付配置不可用")
}
@@ -200,8 +201,11 @@ func (s *RecoverOnlinePaymentService) closePending(ctx context.Context, payment
})
}
func (s *RecoverOnlinePaymentService) adapter(paymentMethod string) OnlinePaymentPort {
func (s *RecoverOnlinePaymentService) adapter(paymentMethod string, config *model.WechatConfig) OnlinePaymentPort {
if paymentMethod == constants.RechargeMethodWechat {
if config != nil && config.ProviderType == model.ProviderTypeFuiou {
return s.fuiou
}
return s.wechat
}
if paymentMethod == constants.RechargeMethodAlipay {

View File

@@ -57,26 +57,9 @@ func (s *Service) ArchiveIntegrationDate(ctx context.Context, archiveDate time.T
return s.archiveIntegrationDate(ctx, archiveDate, false)
}
// FinalizePreviousIntegrationMonth 复核并终结上一个完整自然月的 Integration Log 归档
func (s *Service) FinalizePreviousIntegrationMonth(ctx context.Context) error {
now := time.Now().In(s.location)
return s.FinalizeIntegrationMonth(ctx, now.AddDate(0, -1, 0))
}
// FinalizeIntegrationMonth 逐日复核指定完整自然月,并为变化内容创建最终 revision。
func (s *Service) FinalizeIntegrationMonth(ctx context.Context, month time.Time) error {
monthStart := time.Date(month.In(s.location).Year(), month.In(s.location).Month(), 1, 0, 0, 0, 0, s.location)
currentMonth := time.Now().In(s.location)
currentMonthStart := time.Date(currentMonth.Year(), currentMonth.Month(), 1, 0, 0, 0, 0, s.location)
if !monthStart.Before(currentMonthStart) {
return fmt.Errorf("只能终结已经结束的 Integration Log 完整自然月")
}
for date := monthStart; date.Before(monthStart.AddDate(0, 1, 0)); date = date.AddDate(0, 0, 1) {
if err := s.archiveIntegrationDate(ctx, date, true); err != nil {
return fmt.Errorf("终结 %s Integration Log 归档失败: %w", date.Format(time.DateOnly), err)
}
}
return nil
// FinalizeIntegrationDate 为指定已结束自然日形成 Integration Log 最终归档版本
func (s *Service) FinalizeIntegrationDate(ctx context.Context, archiveDate time.Time) error {
return s.archiveIntegrationDate(ctx, archiveDate, true)
}
func (s *Service) archiveIntegrationDate(ctx context.Context, archiveDate time.Time, final bool) error {
@@ -94,22 +77,20 @@ func (s *Service) archiveIntegrationDate(ctx context.Context, archiveDate time.T
if err != nil {
return err
}
if final && run.Status == constants.ArchiveStatusSuccess && run.IsFinal && run.CleanedAt != nil {
terminalCount, countErr := s.integrationTerminalCount(ctx, start, end)
if countErr != nil {
return countErr
}
if terminalCount == 0 {
return nil
}
}
file, err := s.buildIntegrationArchiveFile(ctx, start, end)
if err != nil {
return err
}
defer os.Remove(file.path)
if final {
pending, pendingErr := s.integrationPendingCount(ctx, start, end)
if pendingErr != nil {
return pendingErr
}
if pending > 0 {
_ = s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).Where("id = ?", run.ID).
Updates(map[string]any{"is_final": false, "error_summary": "存在 pending Integration Log无法形成最终归档", "updated_at": time.Now()}).Error
return fmt.Errorf("仍有 %d 条 pending Integration Log无法形成最终归档", pending)
}
}
if run.Status == constants.ArchiveStatusSuccess && run.RecordCount == file.recordCount && run.SHA256 == file.sha256 {
valid, validateErr := s.validateIntegrationRun(ctx, run)
if validateErr == nil && valid && (!final || run.IsFinal) {
@@ -165,7 +146,8 @@ func (s *Service) acquireIntegrationRun(ctx context.Context, run *model.LogArchi
Where("id = ? AND (status <> ? OR updated_at < ?)", run.ID, constants.ArchiveStatusRunning, now.Add(-3*time.Hour)).
Updates(map[string]any{
"status": constants.ArchiveStatusRunning, "revision": revision, "is_final": false,
"attempt_count": gorm.Expr("attempt_count + 1"), "error_summary": "", "completed_at": nil, "updated_at": now,
"attempt_count": gorm.Expr("attempt_count + 1"), "error_summary": "", "completed_at": nil,
"cleanup_started_at": nil, "cleaned_at": nil, "updated_at": now,
})
if result.Error != nil {
return false, fmt.Errorf("锁定 Integration Log 归档任务失败: %w", result.Error)
@@ -320,16 +302,6 @@ func (s *Service) integrationRecordCount(ctx context.Context, start, end time.Ti
return count, nil
}
func (s *Service) integrationPendingCount(ctx context.Context, start, end time.Time) (int64, error) {
var count int64
if err := s.db.WithContext(ctx).Model(&model.IntegrationLog{}).
Where("created_at >= ? AND created_at < ? AND result = ?", start, end, constants.IntegrationResultPending).
Count(&count).Error; err != nil {
return 0, fmt.Errorf("统计 pending Integration Log 失败: %w", err)
}
return count, nil
}
func integrationObjectKeys(date time.Time, revision int) (string, string) {
prefix := fmt.Sprintf("audit-archive/v1/%04d/%02d/%02d", date.Year(), date.Month(), date.Day())
name := fmt.Sprintf("integration-logs-%s-r%d", date.Format(time.DateOnly), revision)

View File

@@ -7,11 +7,9 @@ import (
"io"
"os"
"strconv"
"strings"
"time"
"github.com/bytedance/sonic"
"gorm.io/gorm"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/pkg/constants"
@@ -19,25 +17,18 @@ import (
const maxManifestBytes = 1024 * 1024
// RetentionAudit 描述月度物理清理的统一审计事实
// RetentionAudit 保留系统审计 Writer 的输入兼容类型
type RetentionAudit struct {
EventID string
Month string
Summary string
Result string
ErrorSummary string
RangeStart time.Time
RangeEnd time.Time
EventCount int64
ResourceCount int64
IntegrationCount int64
EventID, Month, Summary, Result, ErrorSummary string
RangeStart, RangeEnd time.Time
EventCount, ResourceCount, IntegrationCount int64
ManifestKeys []string
DurationMS int64
}
// RetentionResult 是月度留存清理的结构化执行结果。
// RetentionResult 是单日留存处理结果。
type RetentionResult struct {
Month string
ArchiveDate string
EventCount int64
ResourceCount int64
IntegrationCount int64
@@ -46,191 +37,178 @@ type RetentionResult struct {
Duration time.Duration
}
type retentionRuns struct {
audit []*model.LogArchiveRun
integration []*model.LogArchiveRun
// RetentionBlockedError 描述阻断日期推进的安全上下文。
type RetentionBlockedError struct {
ArchiveDate string
Source string
Err error
}
// CleanupPreviousMonth 校验并物理清理上一个完整自然月的在线审计日志。
func (s *Service) CleanupPreviousMonth(ctx context.Context) (RetentionResult, error) {
now := time.Now().In(s.location)
return s.CleanupMonth(ctx, now.AddDate(0, -1, 0))
func (e *RetentionBlockedError) Error() string {
return e.ArchiveDate + " " + e.Source + ": " + e.Err.Error()
}
func (e *RetentionBlockedError) Unwrap() error { return e.Err }
// ValidatePreviousMonth 只读校验上一个完整自然月的归档与清理门禁
func (s *Service) ValidatePreviousMonth(ctx context.Context) (RetentionResult, error) {
now := time.Now().In(s.location)
return s.ValidateMonth(ctx, now.AddDate(0, -1, 0))
}
// ValidateMonth 只读校验指定完整自然月,不写清理断点且不删除在线数据。
func (s *Service) ValidateMonth(ctx context.Context, month time.Time) (result RetentionResult, err error) {
// RetainPendingDays 从最早仍在线的日期连续处理至昨天。cleanup 为 false 时只读校验
func (s *Service) RetainPendingDays(ctx context.Context, cleanup bool) ([]RetentionResult, error) {
if s.db == nil || s.store == nil {
return result, fmt.Errorf("日志留存演练数据库或对象存储未配置")
return nil, fmt.Errorf("日志留存数据库或对象存储未配置")
}
start, end, err := s.retentionMonthRange(month)
start, end, err := s.pendingRetentionRange(ctx)
if err != nil {
return result, err
return nil, err
}
results := make([]RetentionResult, 0)
var auditBlocked, integrationBlocked *RetentionBlockedError
for date := start; date.Before(end); date = date.AddDate(0, 0, 1) {
startedAt := time.Now()
result.Month = start.Format("2006-01")
runs, err := s.loadRetentionRuns(ctx, start, end)
result := RetentionResult{ArchiveDate: date.Format(time.DateOnly)}
if auditBlocked == nil {
auditResult, auditErr := s.retainAuditDate(ctx, date, cleanup)
if auditErr != nil {
auditBlocked = &RetentionBlockedError{ArchiveDate: result.ArchiveDate, Source: constants.AuditArchiveSource, Err: auditErr}
} else {
result.EventCount, result.ResourceCount = auditResult.EventCount, auditResult.ResourceCount
result.ManifestKeys = append(result.ManifestKeys, auditResult.ManifestKeys...)
}
}
if integrationBlocked == nil {
integrationResult, integrationErr := s.retainIntegrationDate(ctx, date, cleanup)
if integrationErr != nil {
integrationBlocked = &RetentionBlockedError{ArchiveDate: result.ArchiveDate, Source: constants.IntegrationArchiveSource, Err: integrationErr}
} else {
result.IntegrationCount = integrationResult.IntegrationCount
result.ManifestKeys = append(result.ManifestKeys, integrationResult.ManifestKeys...)
}
}
result.EstimatedBatches = estimatedRetentionBatches(result)
result.Duration = time.Since(startedAt)
if auditBlocked == nil || integrationBlocked == nil {
results = append(results, result)
}
}
if auditBlocked != nil {
return results, auditBlocked
}
if integrationBlocked != nil {
return results, integrationBlocked
}
return results, nil
}
// RetainDate 校验并按需清理指定已结束自然日,供受控演练使用。
func (s *Service) RetainDate(ctx context.Context, date time.Time, cleanup bool) (RetentionResult, error) {
return s.retainDate(ctx, date, cleanup)
}
func (s *Service) pendingRetentionRange(ctx context.Context) (time.Time, time.Time, error) {
today := time.Now().In(s.location)
end := time.Date(today.Year(), today.Month(), today.Day(), 0, 0, 0, 0, s.location)
var earliest *time.Time
for _, item := range []struct{ table, column, condition string }{
{"tb_audit_event", "created_at", ""},
{"tb_integration_log", "created_at", "result <> 'pending'"},
} {
query := s.db.WithContext(ctx).Table(item.table).Select("MIN(" + item.column + ")")
if item.condition != "" {
query = query.Where(item.condition)
}
var value *time.Time
if err := query.Scan(&value).Error; err != nil {
return time.Time{}, time.Time{}, fmt.Errorf("查询日留存起点失败: %w", err)
}
if value != nil && (earliest == nil || value.Before(*earliest)) {
local := value.In(s.location)
earliest = &local
}
}
if earliest == nil {
return end, end, nil
}
start := time.Date(earliest.Year(), earliest.Month(), earliest.Day(), 0, 0, 0, 0, s.location)
return start, end, nil
}
func (s *Service) retainDate(ctx context.Context, date time.Time, cleanup bool) (RetentionResult, error) {
startedAt := time.Now()
auditResult, err := s.retainAuditDate(ctx, date, cleanup)
if err != nil {
return result, err
return RetentionResult{}, err
}
if err := s.validateRetentionRuns(ctx, start, end, runs); err != nil {
return result, err
integrationResult, err := s.retainIntegrationDate(ctx, date, cleanup)
if err != nil {
return RetentionResult{}, err
}
summarizeRetentionRuns(runs, &result)
result := RetentionResult{ArchiveDate: date.Format(time.DateOnly), EventCount: auditResult.EventCount, ResourceCount: auditResult.ResourceCount, IntegrationCount: integrationResult.IntegrationCount, ManifestKeys: append(auditResult.ManifestKeys, integrationResult.ManifestKeys...)}
result.EstimatedBatches = estimatedRetentionBatches(result)
result.Duration = time.Since(startedAt)
return result, nil
}
// CleanupMonth 校验归档硬门禁后按固定顺序物理清理指定完整自然月。
func (s *Service) CleanupMonth(ctx context.Context, month time.Time) (result RetentionResult, cleanupErr error) {
if s.db == nil || s.store == nil || s.audit == nil {
return result, fmt.Errorf("日志留存清理数据库、对象存储或审计 Writer 未配置")
}
start, end, err := s.retentionMonthRange(month)
func (s *Service) retainAuditDate(ctx context.Context, date time.Time, cleanup bool) (RetentionResult, error) {
run, err := s.retentionRun(ctx, date, constants.AuditArchiveSource)
if err != nil {
return RetentionResult{}, err
}
if run == nil || run.Status != constants.ArchiveStatusSuccess {
if err := s.ArchiveDate(ctx, date); err != nil {
return RetentionResult{}, fmt.Errorf("Audit 归档失败: %w", err)
}
run, err = s.retentionRun(ctx, date, constants.AuditArchiveSource)
if err != nil {
return RetentionResult{}, err
}
}
if err := s.validateAuditRetentionDay(ctx, date, run); err != nil {
return RetentionResult{}, fmt.Errorf("Audit 完整性校验失败: %w", err)
}
result := RetentionResult{ArchiveDate: date.Format(time.DateOnly), EventCount: run.EventCount, ResourceCount: run.ResourceCount, ManifestKeys: []string{run.ManifestKey}}
if cleanup {
if err := s.cleanupAuditDate(ctx, date, run); err != nil {
return result, err
}
startedAt := time.Now()
result.Month = start.Format("2006-01")
cleanupErr = s.executeRetention(ctx, start, end, &result)
result.Duration = time.Since(startedAt)
if auditErr := s.recordRetentionAudit(ctx, start, end, result, cleanupErr); auditErr != nil {
if cleanupErr != nil {
return result, fmt.Errorf("%w记录留存清理失败审计失败: %v", cleanupErr, auditErr)
}
return result, fmt.Errorf("记录留存清理成功审计失败: %w", auditErr)
}
return result, cleanupErr
return result, nil
}
func (s *Service) retentionMonthRange(month time.Time) (time.Time, time.Time, error) {
start := time.Date(month.In(s.location).Year(), month.In(s.location).Month(), 1, 0, 0, 0, 0, s.location)
end := start.AddDate(0, 1, 0)
now := time.Now().In(s.location)
currentMonth := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, s.location)
if !end.Before(currentMonth) && !end.Equal(currentMonth) {
return time.Time{}, time.Time{}, fmt.Errorf("只能清理已经结束的完整自然月")
}
return start, end, nil
}
func (s *Service) executeRetention(ctx context.Context, start, end time.Time, result *RetentionResult) error {
started, err := s.retentionCleanupStarted(ctx, start, end)
func (s *Service) retainIntegrationDate(ctx context.Context, date time.Time, cleanup bool) (RetentionResult, error) {
run, err := s.retentionRun(ctx, date, constants.IntegrationArchiveSource)
if err != nil {
return err
return RetentionResult{}, err
}
if !started {
lastDay := end.AddDate(0, 0, -1)
if err := s.ArchiveDate(ctx, lastDay); err != nil {
return fmt.Errorf("完成上月最后一天 Audit 归档失败: %w", err)
}
if err := s.ArchiveIntegrationDate(ctx, lastDay); err != nil {
return fmt.Errorf("完成上月最后一天 Integration Log 归档失败: %w", err)
}
if err := s.FinalizeIntegrationMonth(ctx, start); err != nil {
return err
if run == nil || run.Status != constants.ArchiveStatusSuccess {
if err := s.ArchiveIntegrationDate(ctx, date); err != nil {
return RetentionResult{}, fmt.Errorf("Integration Log 归档失败: %w", err)
}
}
runs, err := s.loadRetentionRuns(ctx, start, end)
if err := s.FinalizeIntegrationDate(ctx, date); err != nil {
return RetentionResult{}, fmt.Errorf("Integration Log 最终归档失败: %w", err)
}
run, err = s.retentionRun(ctx, date, constants.IntegrationArchiveSource)
if err != nil {
return err
return RetentionResult{}, err
}
if err := s.validateRetentionRuns(ctx, start, end, runs); err != nil {
return err
if err := s.validateIntegrationRetentionDay(ctx, date, run); err != nil {
return RetentionResult{}, fmt.Errorf("Integration Log 完整性校验失败: %w", err)
}
summarizeRetentionRuns(runs, result)
if err := s.cleanupAuditMonth(ctx, start, end, runs.audit); err != nil {
return err
result := RetentionResult{ArchiveDate: date.Format(time.DateOnly), IntegrationCount: run.RecordCount, ManifestKeys: []string{run.ManifestKey}}
if cleanup {
if err := s.cleanupIntegrationDate(ctx, date, run); err != nil {
return result, err
}
return s.cleanupIntegrationMonth(ctx, start, end, runs.integration)
}
return result, nil
}
func (s *Service) retentionCleanupStarted(ctx context.Context, start, end time.Time) (bool, error) {
var count int64
err := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).
Where("archive_date >= ? AND archive_date < ? AND instance_id = ? AND cleanup_started_at IS NOT NULL",
start.Format(time.DateOnly), end.Format(time.DateOnly), s.instanceID).
Count(&count).Error
if err != nil {
return false, fmt.Errorf("读取月度清理断点失败: %w", err)
func (s *Service) retentionRun(ctx context.Context, date time.Time, source string) (*model.LogArchiveRun, error) {
var runs []model.LogArchiveRun
if err := s.db.WithContext(ctx).Where("source = ? AND archive_date = ? AND instance_id = ?", source, date.Format(time.DateOnly), s.instanceID).Find(&runs).Error; err != nil {
return nil, fmt.Errorf("读取日归档账本失败: %w", err)
}
return count > 0, nil
}
func (s *Service) loadRetentionRuns(ctx context.Context, start, end time.Time) (retentionRuns, error) {
var rows []model.LogArchiveRun
err := s.db.WithContext(ctx).Where(
"source IN ? AND archive_date >= ? AND archive_date < ? AND instance_id = ?",
[]string{constants.AuditArchiveSource, constants.IntegrationArchiveSource}, start.Format(time.DateOnly), end.Format(time.DateOnly), s.instanceID,
).Order("archive_date ASC, source ASC").Find(&rows).Error
if err != nil {
return retentionRuns{}, fmt.Errorf("读取月度归档账本失败: %w", err)
if len(runs) == 0 {
return nil, nil
}
days := int(end.Sub(start).Hours() / 24)
if len(rows) != days*2 {
return retentionRuns{}, fmt.Errorf("月度归档账本缺日:期望 %d 条,实际 %d 条", days*2, len(rows))
}
runs := retentionRuns{audit: make([]*model.LogArchiveRun, 0, days), integration: make([]*model.LogArchiveRun, 0, days)}
for index := range rows {
run := &rows[index]
switch run.Source {
case constants.AuditArchiveSource:
runs.audit = append(runs.audit, run)
case constants.IntegrationArchiveSource:
runs.integration = append(runs.integration, run)
}
}
if len(runs.audit) != days || len(runs.integration) != days {
return retentionRuns{}, fmt.Errorf("月度 Audit 或 Integration 归档账本不完整")
}
return runs, nil
}
func (s *Service) validateRetentionRuns(ctx context.Context, start, end time.Time, runs retentionRuns) error {
if err := validateCleanupLedgerState(runs.audit); err != nil {
return fmt.Errorf("Audit 清理断点非法: %w", err)
}
if err := validateCleanupLedgerState(runs.integration); err != nil {
return fmt.Errorf("Integration 清理断点非法: %w", err)
}
for index := range runs.audit {
date := start.AddDate(0, 0, index)
if err := s.validateAuditRetentionDay(ctx, date, runs.audit[index]); err != nil {
return fmt.Errorf("%s Audit 清理门禁失败: %w", date.Format(time.DateOnly), err)
}
if err := s.validateIntegrationRetentionDay(ctx, date, runs.integration[index]); err != nil {
return fmt.Errorf("%s Integration 清理门禁失败: %w", date.Format(time.DateOnly), err)
}
}
return nil
}
func validateCleanupLedgerState(runs []*model.LogArchiveRun) error {
started, cleaned := 0, 0
for _, run := range runs {
if run.CleanupStartedAt != nil {
started++
}
if run.CleanedAt != nil {
cleaned++
}
}
if started != 0 && started != len(runs) {
return fmt.Errorf("清理开始断点不是整月原子状态")
}
if cleaned != 0 && cleaned != len(runs) {
return fmt.Errorf("清理完成断点不是整月原子状态")
}
if cleaned > 0 && started == 0 {
return fmt.Errorf("清理完成但缺少开始断点")
}
return nil
return &runs[0], nil
}
func (s *Service) validateAuditRetentionDay(ctx context.Context, date time.Time, run *model.LogArchiveRun) error {
@@ -258,18 +236,17 @@ func (s *Service) validateIntegrationRetentionDay(ctx context.Context, date time
if err != nil {
return err
}
if run.CleanedAt != nil {
if count != 0 {
return fmt.Errorf("已标记清理完成但数据库仍有 %d 条记录", count)
terminalCount, err := s.integrationTerminalCount(ctx, run.RangeStart, run.RangeEnd)
if err != nil {
return err
}
return nil
if run.CleanedAt != nil && terminalCount != 0 {
return fmt.Errorf("已标记清理完成但数据库仍有 %d 条终态记录", terminalCount)
}
if run.CleanupStartedAt != nil {
if count > run.RecordCount {
if run.CleanupStartedAt != nil && count > run.RecordCount {
return fmt.Errorf("续跑窗口记录数超过最终归档数量")
}
return nil
}
if terminalCount > 0 && run.CleanupStartedAt == nil {
file, err := s.buildIntegrationArchiveFile(ctx, run.RangeStart, run.RangeEnd)
if err != nil {
return err
@@ -278,6 +255,7 @@ func (s *Service) validateIntegrationRetentionDay(ctx context.Context, date time
if file.recordCount != run.RecordCount || file.sha256 != run.SHA256 {
return fmt.Errorf("数据库当前 Integration 内容与最终 revision 不一致")
}
}
return nil
}
@@ -285,8 +263,7 @@ func validateRunBase(run *model.LogArchiveRun, date time.Time, schema string, fi
if run.Status != constants.ArchiveStatusSuccess || run.SchemaVersion != schema {
return fmt.Errorf("归档状态或 schema version 不符合清理要求")
}
if run.ArchiveDate.Format(time.DateOnly) != date.Format(time.DateOnly) ||
!run.RangeStart.Equal(date) || !run.RangeEnd.Equal(date.AddDate(0, 0, 1)) {
if run.ArchiveDate.Format(time.DateOnly) != date.Format(time.DateOnly) || !run.RangeStart.Equal(date) || !run.RangeEnd.Equal(date.AddDate(0, 0, 1)) {
return fmt.Errorf("归档日期或半开时间范围不一致")
}
if final && !run.IsFinal {
@@ -297,21 +274,14 @@ func validateRunBase(run *model.LogArchiveRun, date time.Time, schema string, fi
}
return nil
}
func validateRemainingCounts(run *model.LogArchiveRun, events, resources int64) error {
if run.CleanedAt != nil {
if events != 0 || resources != 0 {
if run.CleanedAt != nil && (events != 0 || resources != 0) {
return fmt.Errorf("已标记清理完成但数据库仍有事件或资源")
}
return nil
}
if run.CleanupStartedAt != nil {
if events > run.EventCount || resources > run.ResourceCount {
if run.CleanupStartedAt != nil && (events > run.EventCount || resources > run.ResourceCount) {
return fmt.Errorf("续跑窗口数量超过已归档数量")
}
return nil
}
if events != run.EventCount || resources != run.ResourceCount {
if run.CleanupStartedAt == nil && (events != run.EventCount || resources != run.ResourceCount) {
return fmt.Errorf("数据库事件或资源数量与 manifest 不一致")
}
return nil
@@ -322,43 +292,27 @@ func (s *Service) validateAuditManifest(ctx context.Context, run *model.LogArchi
if err := s.readManifest(ctx, run.ManifestKey, &manifest); err != nil {
return err
}
if manifest.Source != run.Source || manifest.SchemaVersion != run.SchemaVersion || manifest.Status != constants.ArchiveStatusSuccess ||
manifest.ArchiveDate != run.ArchiveDate.Format(time.DateOnly) || manifest.Timezone != constants.AuditArchiveTimezone ||
manifest.InstanceID != run.InstanceID || !manifest.RangeStart.Equal(run.RangeStart) || !manifest.RangeEnd.Equal(run.RangeEnd) ||
manifest.EventCount != run.EventCount || manifest.ResourceCount != run.ResourceCount ||
manifest.CompressedBytes != run.CompressedBytes || manifest.ObjectKey != run.ObjectKey ||
manifest.SHA256 != run.SHA256 || manifest.Revision != run.Revision {
if manifest.Source != run.Source || manifest.SchemaVersion != run.SchemaVersion || manifest.Status != constants.ArchiveStatusSuccess || manifest.ArchiveDate != run.ArchiveDate.Format(time.DateOnly) || manifest.Timezone != constants.AuditArchiveTimezone || manifest.InstanceID != run.InstanceID || !manifest.RangeStart.Equal(run.RangeStart) || !manifest.RangeEnd.Equal(run.RangeEnd) || manifest.EventCount != run.EventCount || manifest.ResourceCount != run.ResourceCount || manifest.CompressedBytes != run.CompressedBytes || manifest.ObjectKey != run.ObjectKey || manifest.SHA256 != run.SHA256 || manifest.Revision != run.Revision {
return fmt.Errorf("Audit manifest 与 ledger 不一致")
}
if err := s.verifyObject(ctx, run.ManifestKey, -1, map[string]string{
"source": constants.AuditArchiveSource, "data-sha256": run.SHA256, "revision": strconv.Itoa(run.Revision),
}); err != nil {
if err := s.verifyObject(ctx, run.ManifestKey, -1, map[string]string{"source": constants.AuditArchiveSource, "data-sha256": run.SHA256, "revision": strconv.Itoa(run.Revision)}); err != nil {
return err
}
return s.verifyRetentionObject(ctx, run, false)
}
func (s *Service) validateIntegrationManifest(ctx context.Context, run *model.LogArchiveRun) error {
var manifest integrationArchiveManifest
if err := s.readManifest(ctx, run.ManifestKey, &manifest); err != nil {
return err
}
if manifest.Source != run.Source || manifest.SchemaVersion != run.SchemaVersion || manifest.Status != constants.ArchiveStatusSuccess || !manifest.Final ||
manifest.ArchiveDate != run.ArchiveDate.Format(time.DateOnly) || manifest.Timezone != constants.AuditArchiveTimezone ||
manifest.InstanceID != run.InstanceID || !manifest.RangeStart.Equal(run.RangeStart) || !manifest.RangeEnd.Equal(run.RangeEnd) ||
manifest.RecordCount != run.RecordCount || manifest.CompressedBytes != run.CompressedBytes ||
manifest.ObjectKey != run.ObjectKey || manifest.SHA256 != run.SHA256 || manifest.Revision != run.Revision {
if manifest.Source != run.Source || manifest.SchemaVersion != run.SchemaVersion || manifest.Status != constants.ArchiveStatusSuccess || !manifest.Final || manifest.ArchiveDate != run.ArchiveDate.Format(time.DateOnly) || manifest.Timezone != constants.AuditArchiveTimezone || manifest.InstanceID != run.InstanceID || !manifest.RangeStart.Equal(run.RangeStart) || !manifest.RangeEnd.Equal(run.RangeEnd) || manifest.RecordCount != run.RecordCount || manifest.CompressedBytes != run.CompressedBytes || manifest.ObjectKey != run.ObjectKey || manifest.SHA256 != run.SHA256 || manifest.Revision != run.Revision {
return fmt.Errorf("Integration manifest 与最终 ledger 不一致")
}
if err := s.verifyObject(ctx, run.ManifestKey, -1, map[string]string{
"source": constants.IntegrationArchiveSource, "data-sha256": run.SHA256,
"revision": strconv.Itoa(run.Revision), "final": "true",
}); err != nil {
if err := s.verifyObject(ctx, run.ManifestKey, -1, map[string]string{"source": constants.IntegrationArchiveSource, "data-sha256": run.SHA256, "revision": strconv.Itoa(run.Revision), "final": "true"}); err != nil {
return err
}
return s.verifyRetentionObject(ctx, run, true)
}
func (s *Service) readManifest(ctx context.Context, key string, target any) error {
object, err := s.store.Stat(ctx, key)
if err != nil {
@@ -384,13 +338,8 @@ func (s *Service) readManifest(ctx context.Context, key string, target any) erro
}
return nil
}
func (s *Service) verifyRetentionObject(ctx context.Context, run *model.LogArchiveRun, final bool) error {
metadata := map[string]string{
"schema-version": run.SchemaVersion, "source": run.Source,
"archive-date": run.RangeStart.Format(time.DateOnly), "timezone": constants.AuditArchiveTimezone,
"sha256": run.SHA256, "revision": strconv.Itoa(run.Revision),
}
metadata := map[string]string{"schema-version": run.SchemaVersion, "source": run.Source, "archive-date": run.RangeStart.Format(time.DateOnly), "timezone": constants.AuditArchiveTimezone, "sha256": run.SHA256, "revision": strconv.Itoa(run.Revision)}
if run.Source == constants.AuditArchiveSource {
metadata["event-count"] = strconv.FormatInt(run.EventCount, 10)
metadata["resource-count"] = strconv.FormatInt(run.ResourceCount, 10)
@@ -419,53 +368,39 @@ func (s *Service) verifyRetentionObject(ctx context.Context, run *model.LogArchi
}
return nil
}
func summarizeRetentionRuns(runs retentionRuns, result *RetentionResult) {
result.ManifestKeys = make([]string, 0, len(runs.audit)+len(runs.integration))
for _, run := range runs.audit {
result.EventCount += run.EventCount
result.ResourceCount += run.ResourceCount
result.ManifestKeys = append(result.ManifestKeys, run.ManifestKey)
}
for _, run := range runs.integration {
result.IntegrationCount += run.RecordCount
result.ManifestKeys = append(result.ManifestKeys, run.ManifestKey)
}
}
func estimatedRetentionBatches(result RetentionResult) int64 {
batchSize := int64(constants.AuditRetentionDeleteBatchSize)
return (result.EventCount+batchSize-1)/batchSize +
(result.ResourceCount+batchSize-1)/batchSize +
(result.IntegrationCount+batchSize-1)/batchSize
batch := int64(constants.AuditRetentionDeleteBatchSize)
return (result.EventCount+batch-1)/batch + (result.ResourceCount+batch-1)/batch + (result.IntegrationCount+batch-1)/batch
}
func (s *Service) cleanupAuditMonth(ctx context.Context, start, end time.Time, runs []*model.LogArchiveRun) error {
if allRunsCleaned(runs) {
func (s *Service) cleanupAuditDate(ctx context.Context, date time.Time, run *model.LogArchiveRun) error {
if run.CleanedAt != nil {
return nil
}
if err := s.markCleanupStarted(ctx, constants.AuditArchiveSource, start, end); err != nil {
if err := s.markCleanupStarted(ctx, constants.AuditArchiveSource, date); err != nil {
return err
}
if err := s.deleteAuditResources(ctx, start, end); err != nil {
if err := s.deleteAuditResources(ctx, date, date.AddDate(0, 0, 1)); err != nil {
return err
}
if err := s.deleteAuditEvents(ctx, start, end); err != nil {
if err := s.deleteAuditEvents(ctx, date, date.AddDate(0, 0, 1)); err != nil {
return err
}
return s.markCleaned(ctx, constants.AuditArchiveSource, start, end)
return s.markCleaned(ctx, constants.AuditArchiveSource, date)
}
func (s *Service) cleanupIntegrationMonth(ctx context.Context, start, end time.Time, runs []*model.LogArchiveRun) error {
if allRunsCleaned(runs) {
func (s *Service) cleanupIntegrationDate(ctx context.Context, date time.Time, run *model.LogArchiveRun) error {
terminalCount, err := s.integrationTerminalCount(ctx, date, date.AddDate(0, 0, 1))
if err != nil {
return err
}
if run.CleanedAt != nil && terminalCount == 0 {
return nil
}
if err := s.markCleanupStarted(ctx, constants.IntegrationArchiveSource, start, end); err != nil {
if err := s.markCleanupStarted(ctx, constants.IntegrationArchiveSource, date); err != nil {
return err
}
end := date.AddDate(0, 0, 1)
for {
subquery := s.db.Model(&model.IntegrationLog{}).Select("id").
Where("created_at >= ? AND created_at < ?", start, end).Order("id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
subquery := s.db.Model(&model.IntegrationLog{}).Select("id").Where("created_at >= ? AND created_at < ? AND result <> ?", date, end, constants.IntegrationResultPending).Order("id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
deleted := s.db.WithContext(ctx).Where("id IN (?)", subquery).Delete(&model.IntegrationLog{})
if deleted.Error != nil {
return fmt.Errorf("分批物理删除 Integration Log 失败: %w", deleted.Error)
@@ -474,15 +409,21 @@ func (s *Service) cleanupIntegrationMonth(ctx context.Context, start, end time.T
break
}
}
return s.markCleaned(ctx, constants.IntegrationArchiveSource, start, end)
return s.markCleaned(ctx, constants.IntegrationArchiveSource, date)
}
func (s *Service) integrationTerminalCount(ctx context.Context, start, end time.Time) (int64, error) {
var count int64
if err := s.db.WithContext(ctx).Model(&model.IntegrationLog{}).
Where("created_at >= ? AND created_at < ? AND result <> ?", start, end, constants.IntegrationResultPending).
Count(&count).Error; err != nil {
return 0, fmt.Errorf("统计终态 Integration Log 失败: %w", err)
}
return count, nil
}
func (s *Service) deleteAuditResources(ctx context.Context, start, end time.Time) error {
for {
subquery := s.db.Model(&model.AuditEventResource{}).Select("tb_audit_event_resource.id").
Joins("JOIN tb_audit_event ON tb_audit_event.id = tb_audit_event_resource.audit_event_id").
Where("tb_audit_event.created_at >= ? AND tb_audit_event.created_at < ?", start, end).
Order("tb_audit_event_resource.id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
subquery := s.db.Model(&model.AuditEventResource{}).Select("tb_audit_event_resource.id").Joins("JOIN tb_audit_event ON tb_audit_event.id = tb_audit_event_resource.audit_event_id").Where("tb_audit_event.created_at >= ? AND tb_audit_event.created_at < ?", start, end).Order("tb_audit_event_resource.id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
deleted := s.db.WithContext(ctx).Where("id IN (?)", subquery).Delete(&model.AuditEventResource{})
if deleted.Error != nil {
return fmt.Errorf("分批物理删除 Audit Event Resource 失败: %w", deleted.Error)
@@ -492,11 +433,9 @@ func (s *Service) deleteAuditResources(ctx context.Context, start, end time.Time
}
}
}
func (s *Service) deleteAuditEvents(ctx context.Context, start, end time.Time) error {
for {
subquery := s.db.Model(&model.AuditEvent{}).Select("id").
Where("created_at >= ? AND created_at < ?", start, end).Order("id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
subquery := s.db.Model(&model.AuditEvent{}).Select("id").Where("created_at >= ? AND created_at < ?", start, end).Order("id ASC").Limit(constants.AuditRetentionDeleteBatchSize)
deleted := s.db.WithContext(ctx).Where("id IN (?)", subquery).Delete(&model.AuditEvent{})
if deleted.Error != nil {
return fmt.Errorf("分批物理删除 Audit Event 失败: %w", deleted.Error)
@@ -506,75 +445,19 @@ func (s *Service) deleteAuditEvents(ctx context.Context, start, end time.Time) e
}
}
}
func (s *Service) markCleanupStarted(ctx context.Context, source string, start, end time.Time) error {
func (s *Service) markCleanupStarted(ctx context.Context, source string, date time.Time) error {
now := time.Now()
result := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).
Where("source = ? AND archive_date >= ? AND archive_date < ? AND instance_id = ? AND cleanup_started_at IS NULL",
source, start.Format(time.DateOnly), end.Format(time.DateOnly), s.instanceID).
Updates(map[string]any{"cleanup_started_at": now, "updated_at": now})
result := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).Where("source = ? AND archive_date = ? AND instance_id = ? AND cleanup_started_at IS NULL", source, date.Format(time.DateOnly), s.instanceID).Updates(map[string]any{"cleanup_started_at": now, "updated_at": now})
if result.Error != nil {
return fmt.Errorf("记录月度清理开始断点失败: %w", result.Error)
}
return s.validateCleanupMarkerCount(ctx, source, start, end, "cleanup_started_at IS NOT NULL", "开始")
}
func (s *Service) markCleaned(ctx context.Context, source string, start, end time.Time) error {
now := time.Now()
result := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).
Where("source = ? AND archive_date >= ? AND archive_date < ? AND instance_id = ? AND cleanup_started_at IS NOT NULL",
source, start.Format(time.DateOnly), end.Format(time.DateOnly), s.instanceID).
Updates(map[string]any{"cleaned_at": now, "updated_at": now})
if result.Error != nil {
return fmt.Errorf("记录月度清理完成断点失败: %w", result.Error)
}
return s.validateCleanupMarkerCount(ctx, source, start, end, "cleaned_at IS NOT NULL", "完成")
}
func (s *Service) validateCleanupMarkerCount(ctx context.Context, source string, start, end time.Time, marker, label string) error {
var count int64
err := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).
Where("source = ? AND archive_date >= ? AND archive_date < ? AND instance_id = ? AND "+marker,
source, start.Format(time.DateOnly), end.Format(time.DateOnly), s.instanceID).
Count(&count).Error
if err != nil {
return fmt.Errorf("复核月度清理%s断点失败: %w", label, err)
}
expected := int64(end.Sub(start).Hours() / 24)
if count != expected {
return fmt.Errorf("月度清理%s断点不完整期望 %d 条,实际 %d 条", label, expected, count)
return fmt.Errorf("记录清理开始断点失败: %w", result.Error)
}
return nil
}
func allRunsCleaned(runs []*model.LogArchiveRun) bool {
return len(runs) > 0 && runs[0].CleanedAt != nil
}
func (s *Service) recordRetentionAudit(ctx context.Context, start, end time.Time, result RetentionResult, cleanupErr error) error {
audit := RetentionAudit{
Month: result.Month, RangeStart: start, RangeEnd: end,
EventCount: result.EventCount, ResourceCount: result.ResourceCount,
IntegrationCount: result.IntegrationCount, ManifestKeys: result.ManifestKeys,
DurationMS: result.Duration.Milliseconds(), Result: constants.AuditResultSuccess,
Summary: "完成已归档在线日志月度物理清理",
EventID: "evt_retention_" + strings.ReplaceAll(result.Month, "-", "_"),
func (s *Service) markCleaned(ctx context.Context, source string, date time.Time) error {
now := time.Now()
result := s.db.WithContext(ctx).Model(&model.LogArchiveRun{}).Where("source = ? AND archive_date = ? AND instance_id = ? AND cleanup_started_at IS NOT NULL", source, date.Format(time.DateOnly), s.instanceID).Updates(map[string]any{"cleaned_at": now, "updated_at": now})
if result.Error != nil {
return fmt.Errorf("记录日清理完成断点失败: %w", result.Error)
}
if cleanupErr != nil {
audit.EventID = ""
audit.Result = constants.AuditResultFailed
audit.Summary = "已归档在线日志月度物理清理失败"
audit.ErrorSummary = truncateRetentionError(cleanupErr)
}
return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
return s.audit.WriteRetentionCleanup(ctx, tx, audit)
})
}
func truncateRetentionError(err error) string {
value := []rune(err.Error())
if len(value) > 500 {
value = value[:500]
}
return string(value)
return nil
}

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

@@ -288,6 +288,7 @@ func initServices(s *stores, deps *Dependencies) *services {
deps.DB,
paymentInfra.NewWechatWebAdapter(wechat.NewRedisCache(deps.Redis), paymentIntegration, deps.Logger),
paymentInfra.NewAlipayWapAdapter(paymentIntegration, deps.Logger),
paymentInfra.NewFuiouScanAdapter(paymentIntegration, deps.Logger),
auditWriter,
)
agentRechargePaymentConfirm := agentrechargeApp.NewConfirmOnlinePaymentService(

View File

@@ -3,6 +3,7 @@ package agentrecharge
import (
"strings"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/errors"
)
@@ -51,8 +52,11 @@ func ValidatePaymentConfirmation(facts PaymentConfirmationFacts) (bool, error) {
return false, errors.New(errors.CodeConflict, "支付单与代理充值单关联不一致")
}
method := strings.TrimSpace(facts.PaymentMethod)
if method == "" || method != strings.TrimSpace(facts.RechargePaymentMethod) ||
method != strings.TrimSpace(facts.RechargePaymentChannel) {
if method == "" || method != strings.TrimSpace(facts.RechargePaymentMethod) {
return false, errors.New(errors.CodeConflict, "支付渠道与代理充值单不一致")
}
channel := strings.TrimSpace(facts.RechargePaymentChannel)
if channel != method && !(method == constants.RechargeMethodWechat && channel == model.ProviderTypeFuiou) {
return false, errors.New(errors.CodeConflict, "支付渠道与代理充值单不一致")
}
identity := strings.TrimSpace(facts.MerchantIdentity)

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

@@ -283,8 +283,13 @@ func (h *PaymentHandler) confirmAgentRechargePayment(ctx context.Context, callba
correlationID := callback.PaymentNo
ctx = auditcontext.With(ctx, auditcontext.Context{CorrelationID: correlationID})
linkage := auditcontext.From(ctx)
// 富友本质是微信支付上游通道,回调渠道归一化为业务方式 wechat 后再进入确认用例。
paymentMethod := callback.PaymentMethod
if paymentMethod == model.ProviderTypeFuiou {
paymentMethod = constants.RechargeMethodWechat
}
result, confirmErr := h.agentPaymentConfirm.Execute(ctx, agentrechargeApp.ConfirmOnlinePaymentCommand{
PaymentNo: callback.PaymentNo, PaymentMethod: callback.PaymentMethod, ConfigID: callback.ConfigID,
PaymentNo: callback.PaymentNo, PaymentMethod: paymentMethod, ConfigID: callback.ConfigID,
MerchantIdentity: callback.MerchantIdentity, ThirdPartyTradeNo: callback.TransactionID,
Amount: callback.Amount, PaidAt: callback.PaidAt, RequestID: linkage.RequestID,
CorrelationID: correlationID, ParentEventID: linkage.ParentEventID,

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

@@ -0,0 +1,172 @@
package payment
import (
"context"
"strconv"
"strings"
"time"
"go.uber.org/zap"
agentrecharge "github.com/break/junhong_cmp_fiber/internal/application/agentrecharge"
"github.com/break/junhong_cmp_fiber/internal/infrastructure/integrationlog"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/pkg/constants"
apperrors "github.com/break/junhong_cmp_fiber/pkg/errors"
"github.com/break/junhong_cmp_fiber/pkg/fuiou"
)
// FuiouScanAdapter 按富友主扫统一下单生成微信扫码支付链接并主动查单。
type FuiouScanAdapter struct {
integration *integrationlog.Repository
logger *zap.Logger
}
// NewFuiouScanAdapter 创建富友扫码支付适配器。
func NewFuiouScanAdapter(integration *integrationlog.Repository, logger *zap.Logger) *FuiouScanAdapter {
return &FuiouScanAdapter{integration: integration, logger: logger}
}
// Available 判断富友主扫下单、验签与查单所需的配置是否完整。
func (a *FuiouScanAdapter) Available(config *model.WechatConfig) bool {
return fuiouConfigComplete(config, true)
}
// CreatePaymentURL 调用富友主扫统一下单并返回二维码链接。
func (a *FuiouScanAdapter) CreatePaymentURL(ctx context.Context, request agentrecharge.OnlinePaymentRequest) (agentrecharge.OnlinePaymentResult, error) {
attempt, err := a.startAttempt(ctx, request, constants.IntegrationOperationPaymentPreCreate, request.Config.ID, request.Amount)
if err != nil {
return agentrecharge.OnlinePaymentResult{}, err
}
startedAt := time.Now()
client, err := a.newClient(request.Config)
if err != nil {
return agentrecharge.OnlinePaymentResult{}, a.completeUnknown(ctx, attempt.IntegrationID, startedAt, err)
}
expireMinutes := int(time.Until(request.ExpireAt).Minutes())
if expireMinutes < 1 {
expireMinutes = 1
}
resp, callErr := client.PreCreate(
request.PaymentNo, strconv.FormatInt(request.Amount, 10), request.Description,
fuiou.GetServerIP(), fuiou.OrderTypeWechat, strconv.Itoa(expireMinutes),
)
if callErr != nil {
return agentrecharge.OnlinePaymentResult{}, a.completeUnknown(ctx, attempt.IntegrationID, startedAt, callErr)
}
if strings.TrimSpace(resp.QrCode) == "" {
return agentrecharge.OnlinePaymentResult{}, a.completeFailed(ctx, attempt.IntegrationID, startedAt, "empty_qr_code", "富友主扫下单未返回二维码链接")
}
if _, err = a.integration.Complete(ctx, attempt.IntegrationID, integrationlog.Completion{
Result: constants.IntegrationResultSuccess, ResponseSummary: map[string]any{"success": true},
DurationMS: time.Since(startedAt).Milliseconds(), StateChanged: true,
}); err != nil {
return agentrecharge.OnlinePaymentResult{}, err
}
return agentrecharge.OnlinePaymentResult{QRContent: resp.QrCode}, nil
}
// Query 调用富友订单查询并按 trans_stat 返回统一查询结果。
func (a *FuiouScanAdapter) Query(ctx context.Context, request agentrecharge.OnlinePaymentRequest) (agentrecharge.OnlinePaymentQueryResult, error) {
attempt, err := a.startAttempt(ctx, request, constants.IntegrationOperationPaymentQuery, request.Config.ID, 0)
if err != nil {
return agentrecharge.OnlinePaymentQueryResult{}, err
}
startedAt := time.Now()
client, err := a.newClient(request.Config)
if err != nil {
return agentrecharge.OnlinePaymentQueryResult{}, a.completeUnknown(ctx, attempt.IntegrationID, startedAt, err)
}
resp, callErr := client.CommonQuery(request.PaymentNo, fuiou.OrderTypeWechat)
if callErr != nil {
return agentrecharge.OnlinePaymentQueryResult{}, a.completeUnknown(ctx, attempt.IntegrationID, startedAt, callErr)
}
result := agentrecharge.OnlinePaymentQueryResult{State: mapFuiouTransStat(resp.TransStat)}
if result.State == agentrecharge.OnlinePaymentStatePaid {
result.ThirdPartyTradeNo = strings.TrimSpace(resp.TransactionId)
result.Amount, _ = strconv.ParseInt(strings.TrimSpace(resp.OrderAmt), 10, 64)
if paidAt, ok := parseWechatPaidAt(strings.TrimSpace(resp.ReservedTxnFinTs)); ok {
result.PaidAt = &paidAt
}
}
providerCode := resp.TransStat
if strings.TrimSpace(providerCode) == "" {
providerCode = "unknown"
}
if _, err = a.integration.Complete(ctx, attempt.IntegrationID, integrationlog.Completion{
Result: constants.IntegrationResultSuccess, ProviderCode: providerCode,
ResponseSummary: map[string]any{"state": result.State, "has_trade_no": result.ThirdPartyTradeNo != ""},
DurationMS: time.Since(startedAt).Milliseconds(),
}); err != nil {
return agentrecharge.OnlinePaymentQueryResult{}, err
}
return result, nil
}
func fuiouConfigComplete(config *model.WechatConfig, requireActive bool) bool {
return config != nil && (!requireActive || config.IsActive) && config.ProviderType == model.ProviderTypeFuiou &&
config.FyInsCd != "" && config.FyMchntCd != "" && config.FyTermID != "" &&
config.FyPrivateKey != "" && config.FyPublicKey != "" && config.FyAPIURL != "" && config.FyNotifyURL != ""
}
func (a *FuiouScanAdapter) newClient(config *model.WechatConfig) (*fuiou.Client, error) {
if !fuiouConfigComplete(config, false) {
return nil, apperrors.New(apperrors.CodeNoPaymentConfig, "富友扫码支付配置不可用")
}
return fuiou.NewClient(
config.FyInsCd, config.FyMchntCd, config.FyTermID, config.FyAPIURL, config.FyNotifyURL,
config.FyPrivateKey, config.FyPublicKey, a.logger,
)
}
func mapFuiouTransStat(transStat string) string {
switch transStat {
case "SUCCESS":
return agentrecharge.OnlinePaymentStatePaid
case "PAYERROR", "CLOSED", "REVOKED":
return agentrecharge.OnlinePaymentStateClosed
case "USERPAYING", "NOTPAY":
return agentrecharge.OnlinePaymentStatePending
default:
return agentrecharge.OnlinePaymentStateUnknown
}
}
func (a *FuiouScanAdapter) startAttempt(ctx context.Context, request agentrecharge.OnlinePaymentRequest, operation string, configID uint, amount int64) (*model.IntegrationLog, error) {
resourceID, resourceKey := strconv.FormatUint(uint64(request.PaymentID), 10), request.PaymentNo
series := "agent-recharge-payment:" + resourceID + ":" + operation
correlationID := request.CorrelationID
return a.integration.Start(ctx, integrationlog.Attempt{
Provider: constants.IntegrationProviderFuiou, Direction: constants.IntegrationDirectionOutbound,
Operation: operation, ResourceType: constants.IntegrationResourceTypeAgentRechargePayment,
ResourceID: &resourceID, ResourceKey: &resourceKey, ExternalID: &resourceKey,
TriggerSeries: &series, CorrelationID: &correlationID,
RequestSummary: map[string]any{"payment_config_id": configID, "amount": amount},
})
}
func (a *FuiouScanAdapter) completeUnknown(ctx context.Context, integrationID string, startedAt time.Time, cause error) error {
if a.logger != nil {
a.logger.Warn("富友支付请求结果未知", zap.String("integration_id", integrationID), zap.Error(cause))
}
_, err := a.integration.Complete(ctx, integrationID, integrationlog.Completion{
Result: constants.IntegrationResultUnknown, ProviderCode: "request_unknown", SafeProviderMessage: "富友支付请求结果未知",
ResponseSummary: map[string]any{"success": false}, DurationMS: time.Since(startedAt).Milliseconds(),
RecoveryStrategy: "使用原支付单号主动查单,确认不存在或关闭后才允许关闭本地支付单",
})
if err != nil {
return err
}
return apperrors.Wrap(apperrors.CodeTimeout, cause, "富友支付请求结果未知")
}
func (a *FuiouScanAdapter) completeFailed(ctx context.Context, integrationID string, startedAt time.Time, providerCode, providerMessage string) error {
_, err := a.integration.Complete(ctx, integrationID, integrationlog.Completion{
Result: constants.IntegrationResultFailed, ProviderCode: providerCode, SafeProviderMessage: providerMessage,
ResponseSummary: map[string]any{"success": false}, DurationMS: time.Since(startedAt).Milliseconds(),
})
if err != nil {
return err
}
return apperrors.New(apperrors.CodeServiceUnavailable, providerMessage)
}

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,8 @@ 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,30 +310,17 @@ 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"))
activationEnqueued, err := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
if err != nil {
// 激活结果未知时不能将其当作无套餐;孤儿扫描和套餐轮询会继续兜底。
return err
}
if activationEnqueued {
return nil
}
h.triggerStopAfterExpiry(ctx, carrierType, carrierID)
return activationErr
}
return nil
@@ -401,8 +374,36 @@ func (h *PackageActivationHandler) getCarrierInfo(pkg *model.PackageUsage) (stri
return "", 0
}
// triggerStopAfterExpiry 套餐过期后异步触发停机检查
// 仅在确认无后续生效套餐时有效CheckAndStopCard 内部有幂等保护,重复调用安全
// activateNextPackage 提交下一个待生效主套餐的异步激活任务。
// 返回 true 表示已找到并成功提交后续套餐false 表示不存在后续套餐。
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) (bool, 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 false, nil
}
return false, err
}
if err := h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue"); err != nil {
return false, err
}
return true, nil
}
// triggerStopAfterExpiry 在明确无后续套餐时异步触发停机检查。
// 后续套餐存在时,停机重评估由激活任务在完成后顺序执行。
func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, carrierType string, carrierID uint) {
if h.stopResumeCallback == nil {
return
@@ -441,6 +442,42 @@ func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, c
}
}
// reconcileCarrierAfterActivation 在套餐激活任务完成后按最新事实重评估停复机。
// 它必须在激活调用返回后执行,不能与激活任务并发读取过期权益快照。
func (h *PackageActivationHandler) reconcileCarrierAfterActivation(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint) error {
if h.stopResumeCallback == nil {
h.logger.Warn("套餐激活后停复机回调未注入,跳过重评估",
zap.Uint("package_usage_id", packageUsageID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID))
return nil
}
if carrierID == 0 {
return errors.New(errors.CodeInvalidParam, "套餐使用记录缺少有效载体")
}
if carrierType == constants.AssetTypeIotCard {
if err := h.stopResumeCallback.CheckAndStopCard(ctx, carrierID); err != nil {
return err
}
return nil
}
if carrierType != constants.AssetTypeDevice {
return errors.New(errors.CodeInvalidParam, "套餐使用记录载体类型无效")
}
bindings, err := h.deviceSimBinding.ListByDeviceID(ctx, carrierID)
if err != nil {
return err
}
for _, binding := range bindings {
if err := h.stopResumeCallback.CheckAndStopCard(ctx, binding.IotCardID); err != nil {
return err
}
}
return nil
}
// enqueueActivationTask 提交套餐激活任务到 Asynq
func (h *PackageActivationHandler) enqueueActivationTask(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint, activationType string) error {
linkage := auditcontext.From(ctx)
@@ -509,14 +546,7 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
return err
}
// 幂等性检查:如果已经是生效状态,跳过
if pkg.Status == constants.PackageUsageStatusActive {
h.logger.Info("套餐已激活,跳过",
zap.Uint("package_usage_id", payload.PackageUsageID))
return nil
}
// 调用 ActivationService 执行激活
// 调用 ActivationService 执行激活。即使套餐已由其他任务激活,仍须重新评估停复机。
if h.activationService != nil {
if err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID); err != nil {
h.logger.Error("套餐激活失败",
@@ -531,8 +561,20 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
return errors.New(errors.CodeInternalError, "激活服务未注入,无法执行套餐激活")
}
h.logger.Info("套餐激活成功",
zap.Uint("package_usage_id", payload.PackageUsageID))
carrierType, carrierID := h.getCarrierInfo(&pkg)
if err := h.reconcileCarrierAfterActivation(ctx, payload.PackageUsageID, carrierType, carrierID); err != nil {
h.logger.Error("套餐激活后停复机重评估失败",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.Error(err))
return err
}
h.logger.Info("套餐激活及停复机重评估完成",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID))
return nil
}

View File

@@ -98,6 +98,20 @@ func (m *PollingQueueManager) Requeue(ctx context.Context, cardID uint, taskType
}).Err()
}
// EnsureQueued 仅在任务当前不在分片队列中时补入任务,不覆盖已有任务的执行时间。
func (m *PollingQueueManager) EnsureQueued(ctx context.Context, cardID uint, taskType string, nextCheckAt time.Time) (bool, error) {
shardID := int(cardID) % m.shardCount
key := constants.RedisPollingShardQueueKey(shardID, taskType)
added, err := m.redis.ZAddArgs(ctx, key, redis.ZAddArgs{
NX: true,
Members: []redis.Z{{
Score: float64(nextCheckAt.Unix()),
Member: fmt.Sprintf("%d", cardID),
}},
}).Result()
return added > 0, err
}
// RemoveFromAllQueues 从所有分片的所有5个队列realname/carddata/package/protect/card_status移除指定卡
// 修复 Bug3旧实现漏掉 protect 队列
func (m *PollingQueueManager) RemoveFromAllQueues(ctx context.Context, cardID uint) error {

View File

@@ -13,75 +13,108 @@ import (
"github.com/break/junhong_cmp_fiber/pkg/errors"
)
// Source 表示受在线留存边界约束的数据源。
type Source string
const (
// SourceAudit 表示统一审计事件。
SourceAudit Source = constants.AuditArchiveSource
// SourceIntegration 表示外部交互日志。
SourceIntegration Source = constants.IntegrationArchiveSource
)
// Info 是查询响应公开的在线留存边界。
type Info struct {
OnlineFrom time.Time `json:"online_from" description:"当前可在线查询的最早时间"`
ArchivedBefore *time.Time `json:"archived_before" description:"早于该时间的数据已归档;尚未清理时为空"`
Timezone string `json:"timezone" description:"留存自然日时区"`
}
// Load 从归档账本读取已完成物理清理的数据边界
// Load 仅公开从最早在线日期开始连续完成物理清理的边界,绝不跨越清理空洞
func Load(ctx context.Context, db *gorm.DB, sources ...Source) (Info, error) {
location, err := time.LoadLocation(constants.AuditArchiveTimezone)
if err != nil {
return Info{}, errors.Wrap(errors.CodeInternalError, err, "加载审计留存时区失败")
}
now := time.Now().In(location)
info := Info{OnlineFrom: time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, location), Timezone: constants.AuditArchiveTimezone}
boundaries := make([]sourceRetention, 0, len(sources))
for _, source := range sources {
boundary, cleaned, err := sourceBoundary(ctx, db, source, location)
boundary, err := sourceBoundary(ctx, db, source, location)
if err != nil {
return Info{}, err
}
if boundary.Before(info.OnlineFrom) && info.ArchivedBefore == nil {
info.OnlineFrom = boundary
boundaries = append(boundaries, boundary)
}
if cleaned && (info.ArchivedBefore == nil || boundary.After(*info.ArchivedBefore)) {
value := boundary
now := time.Now().In(location)
fallback := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, location)
info := Info{OnlineFrom: fallback, Timezone: constants.AuditArchiveTimezone}
if len(boundaries) == 0 {
return info, nil
}
for _, boundary := range boundaries {
if boundary.onlineFrom.Before(info.OnlineFrom) {
info.OnlineFrom = boundary.onlineFrom
}
}
if len(boundaries) == 1 && boundaries[0].cleaned {
value := boundaries[0].onlineFrom
info.ArchivedBefore = &value
info.OnlineFrom = boundary
return info, nil
}
if len(boundaries) > 1 {
common := boundaries[0].onlineFrom
allCleaned := boundaries[0].cleaned
for _, boundary := range boundaries[1:] {
if boundary.onlineFrom.Before(common) {
common = boundary.onlineFrom
}
allCleaned = allCleaned && boundary.cleaned
}
if allCleaned {
info.OnlineFrom = common
info.ArchivedBefore = &common
}
}
return info, nil
}
func sourceBoundary(ctx context.Context, db *gorm.DB, source Source, location *time.Location) (time.Time, bool, error) {
var cleanedEnd sql.NullTime
if err := db.WithContext(ctx).Model(&model.LogArchiveRun{}).
Where("source = ? AND cleaned_at IS NOT NULL", source).
Select("MAX(range_end)").Scan(&cleanedEnd).Error; err != nil {
return time.Time{}, false, errors.Wrap(errors.CodeDatabaseError, err, "查询审计留存清理边界失败")
}
if cleanedEnd.Valid {
return cleanedEnd.Time.In(location), true, nil
}
type sourceRetention struct {
onlineFrom time.Time
cleaned bool
}
var earliest sql.NullTime
table, column := "tb_audit_event", "occurred_at"
func sourceBoundary(ctx context.Context, db *gorm.DB, source Source, location *time.Location) (sourceRetention, error) {
table, column := "tb_audit_event", "created_at"
if source == SourceIntegration {
table, column = "tb_integration_log", "created_at"
}
if err := db.WithContext(ctx).Table(table).Select("MIN(" + column + ")").Scan(&earliest).Error; err != nil {
return time.Time{}, false, errors.Wrap(errors.CodeDatabaseError, err, "查询审计在线数据边界失败")
var earliestOnline, earliestLedger sql.NullTime
if err := db.WithContext(ctx).Table(table).Select("MIN(" + column + ")").Scan(&earliestOnline).Error; err != nil {
return sourceRetention{}, errors.Wrap(errors.CodeDatabaseError, err, "查询审计在线数据边界失败")
}
if earliest.Valid {
return earliest.Time.In(location), false, nil
if err := db.WithContext(ctx).Model(&model.LogArchiveRun{}).Where("source = ?", source).Select("MIN(archive_date)").Scan(&earliestLedger).Error; err != nil {
return sourceRetention{}, errors.Wrap(errors.CodeDatabaseError, err, "查询审计留存账本边界失败")
}
if !earliestOnline.Valid && !earliestLedger.Valid {
now := time.Now().In(location)
return time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, location), false, nil
return sourceRetention{onlineFrom: time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, location)}, nil
}
start := earliestLedger.Time.In(location)
if earliestOnline.Valid && (!earliestLedger.Valid || earliestOnline.Time.Before(earliestLedger.Time)) {
start = earliestOnline.Time.In(location)
}
start = time.Date(start.Year(), start.Month(), start.Day(), 0, 0, 0, 0, location)
var rows []model.LogArchiveRun
if err := db.WithContext(ctx).Where("source = ? AND archive_date >= ?", source, start.Format(time.DateOnly)).Order("archive_date ASC").Find(&rows).Error; err != nil {
return sourceRetention{}, errors.Wrap(errors.CodeDatabaseError, err, "查询审计留存清理边界失败")
}
expected := start
for _, row := range rows {
date := row.ArchiveDate
date = time.Date(date.Year(), date.Month(), date.Day(), 0, 0, 0, 0, location)
if !date.Equal(expected) || row.CleanedAt == nil {
break
}
expected = expected.AddDate(0, 0, 1)
}
return sourceRetention{onlineFrom: expected, cleaned: !expected.Equal(start)}, nil
}
// NormalizeRange 将缺省范围收敛到在线窗口,并拒绝归档或跨边界查询。
func NormalizeRange(info Info, from, to *time.Time, maxRange ...time.Duration) (*time.Time, *time.Time, error) {
explicitFrom := from != nil
if from != nil && info.ArchivedBefore != nil && from.Before(info.OnlineFrom) {
@@ -110,7 +143,6 @@ func NormalizeRange(info Info, from, to *time.Time, maxRange ...time.Duration) (
}
return from, to, nil
}
func archivedError(info Info) error {
return errors.NewWithData(errors.CodeAuditDataArchived, map[string]any{"retention": info})
}

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

@@ -334,10 +334,10 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package
return errors.Wrap(errors.CodeRedisError, err, "获取分布式锁失败")
}
if !locked {
s.logger.Warn("套餐激活正在进行中,跳过",
s.logger.Warn("套餐激活正在进行中,等待任务重试",
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID))
return nil
return errors.New(errors.CodePackageActivationConflict)
}
defer s.redis.Del(ctx, lockKey)

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 {
@@ -388,7 +409,7 @@ func (s *Service) appendCompletedNotification(ctx context.Context, tx *gorm.DB,
// refundWalletPayment 处理钱包支付订单的退款回款。
func (s *Service) refundWalletPayment(ctx context.Context, tx *gorm.DB, refund *model.RefundRequest, order *model.Order, amount int64, operatorID uint) error {
if order.PaymentMethod != model.PaymentMethodWallet {
if amount == 0 || order.PaymentMethod != model.PaymentMethodWallet {
return nil
}
@@ -1180,8 +1201,8 @@ func validateRequestedRefundAmountByOrder(requestedRefundAmount int64, order *mo
// validateApprovedRefundAmount 校验审批退款金额不能超过申请金额和订单实收金额。
func validateApprovedRefundAmount(approvedAmount int64, requestedRefundAmount int64, order *model.Order) error {
if approvedAmount <= 0 {
return errors.New(errors.CodeInvalidParam, "审批退款金额必须大于0")
if approvedAmount < 0 {
return errors.New(errors.CodeInvalidParam, "审批退款金额不能小于0")
}
if approvedAmount > requestedRefundAmount {
return errors.New(errors.CodeInvalidParam, "审批退款金额不能大于申请退款金额")

View File

@@ -121,19 +121,13 @@ func (h *AssetPackageBatchOrderHandler) finishBatchOrderTask(ctx context.Context
return err
}
rootID := audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "completed")
var childCount int64
if err := tx.WithContext(ctx).Model(&model.AuditEvent{}).
Where("parent_event_id = ? AND action_code = ?", rootID, constants.AuditActionOrderCreated).
Count(&childCount).Error; err != nil {
return err
}
result := batchAuditResult(int(childCount), failCount)
result := batchAuditResult(successCount, failCount)
return h.auditWriter.WriteTask(ctx, tx, audit.TaskInput{
EventID: rootID, ActionCode: constants.AuditActionAssetPackageBatchOrderTaskCompleted,
Summary: "完成资产套餐批量订购任务", TaskID: taskRecord.ID, TaskNo: taskRecord.TaskNo,
Result: result, CorrelationID: taskRecord.TaskNo,
ParentEventID: audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "created"),
BatchTotal: len(items), SuccessCount: int(childCount), FailCount: failCount,
BatchTotal: len(items), SuccessCount: successCount, FailCount: failCount,
IdentitySnapshot: map[string]any{
"id": taskRecord.ID, "task_no": taskRecord.TaskNo, "file_name": taskRecord.FileName,
"package_id": taskRecord.PackageID, "package_code": taskRecord.PackageCode,

View File

@@ -2,101 +2,47 @@ package task
import (
"context"
stderrors "errors"
"fmt"
"time"
"github.com/bytedance/sonic"
"github.com/hibiken/asynq"
"go.uber.org/zap"
"github.com/break/junhong_cmp_fiber/internal/application/auditarchive"
)
// AuditMonthlyRetentionPayload 是人工补跑月度清理时可选的任务载荷
type AuditMonthlyRetentionPayload struct {
ArchiveMonth string `json:"archive_month"`
}
// AuditMonthlyRetentionHandler 处理归档完整性门禁与上月在线日志物理清理。
type AuditMonthlyRetentionHandler struct {
// AuditRetentionHandler 连续校验并按日物理清理已结束的在线日志
type AuditRetentionHandler struct {
service *auditarchive.Service
logger *zap.Logger
cleanupEnabled bool
}
// NewAuditMonthlyRetentionHandler 创建月度日志留存清理处理器。
func NewAuditMonthlyRetentionHandler(service *auditarchive.Service, logger *zap.Logger, cleanupEnabled bool) *AuditMonthlyRetentionHandler {
return &AuditMonthlyRetentionHandler{service: service, logger: logger, cleanupEnabled: cleanupEnabled}
// NewAuditRetentionHandler 创建日留存处理器。
func NewAuditRetentionHandler(service *auditarchive.Service, logger *zap.Logger, cleanupEnabled bool) *AuditRetentionHandler {
return &AuditRetentionHandler{service: service, logger: logger, cleanupEnabled: cleanupEnabled}
}
// Handle 校验整月归档后按固定顺序分批物理删除 PostgreSQL 在线日志
func (h *AuditMonthlyRetentionHandler) Handle(ctx context.Context, task *asynq.Task) error {
if !h.cleanupEnabled {
return h.handleDryRun(ctx, task)
}
// Handle 在关闭清理开关时只读校验,开启后从最早未完成日连续删除
func (h *AuditRetentionHandler) Handle(ctx context.Context, _ *asynq.Task) error {
if h.service == nil {
return fmt.Errorf("月度日志留存清理服务未配置")
return fmt.Errorf("日志留存服务未配置")
}
startedAt := time.Now()
var result auditarchive.RetentionResult
var err error
if len(task.Payload()) == 0 {
result, err = h.service.CleanupPreviousMonth(ctx)
} else {
var payload AuditMonthlyRetentionPayload
if unmarshalErr := sonic.Unmarshal(task.Payload(), &payload); unmarshalErr != nil {
return fmt.Errorf("解析月度日志留存清理任务载荷失败: %w", unmarshalErr)
}
month, parseErr := parseArchiveMonth(payload.ArchiveMonth)
if parseErr != nil {
return parseErr
}
result, err = h.service.CleanupMonth(ctx, month)
}
fields := []zap.Field{
zap.String("archive_month", result.Month), zap.Int64("audit_event_count", result.EventCount),
zap.Int64("event_resource_count", result.ResourceCount), zap.Int64("integration_log_count", result.IntegrationCount),
zap.Duration("duration", time.Since(startedAt)), zap.Int("manifest_count", len(result.ManifestKeys)),
}
if err != nil {
fields = append(fields, zap.String("severity", "critical"), zap.Error(err))
h.logger.Error("月度日志留存清理失败PostgreSQL 整月清理已阻断或等待断点续跑", fields...)
return err
}
h.logger.Info("月度日志留存清理完成", fields...)
return nil
}
func (h *AuditMonthlyRetentionHandler) handleDryRun(ctx context.Context, task *asynq.Task) error {
if h.service == nil {
return fmt.Errorf("月度日志留存演练服务未配置")
}
var result auditarchive.RetentionResult
var err error
if len(task.Payload()) == 0 {
result, err = h.service.ValidatePreviousMonth(ctx)
} else {
var payload AuditMonthlyRetentionPayload
if unmarshalErr := sonic.Unmarshal(task.Payload(), &payload); unmarshalErr != nil {
return fmt.Errorf("解析月度日志留存演练任务载荷失败: %w", unmarshalErr)
}
month, parseErr := parseArchiveMonth(payload.ArchiveMonth)
if parseErr != nil {
return parseErr
}
result, err = h.service.ValidateMonth(ctx, month)
}
fields := []zap.Field{
zap.Bool("cleanup_enabled", false), zap.String("archive_month", result.Month),
results, err := h.service.RetainPendingDays(ctx, h.cleanupEnabled)
for _, result := range results {
h.logger.Info("日志日留存处理完成",
zap.String("archive_date", result.ArchiveDate), zap.Bool("cleanup_enabled", h.cleanupEnabled),
zap.Int64("audit_event_count", result.EventCount), zap.Int64("event_resource_count", result.ResourceCount),
zap.Int64("integration_log_count", result.IntegrationCount), zap.Int64("estimated_cleanup_batches", result.EstimatedBatches),
zap.Duration("validation_duration", result.Duration), zap.Int("manifest_count", len(result.ManifestKeys)),
zap.Int64("integration_log_count", result.IntegrationCount), zap.Duration("duration", result.Duration))
}
if err != nil {
fields = append(fields, zap.String("severity", "critical"), zap.Error(err))
h.logger.Error("审计日志月度只读演练失败,物理清理保持关闭", fields...)
var blocked *auditarchive.RetentionBlockedError
if stderrors.As(err, &blocked) {
h.logger.Error("日志日留存日期推进已阻断", zap.String("archive_date", blocked.ArchiveDate), zap.String("source", blocked.Source), zap.String("failure_category", "archive_or_validation"), zap.Error(err))
} else {
h.logger.Error("日志日留存日期推进已阻断", zap.String("source", "retention"), zap.String("failure_category", "internal"), zap.Error(err))
}
return err
}
h.logger.Info("审计日志月度只读演练通过,物理清理保持关闭", fields...)
return nil
}

View File

@@ -18,12 +18,7 @@ type IntegrationDailyArchivePayload struct {
ArchiveDate string `json:"archive_date"`
}
// IntegrationMonthlyFinalizePayload 是人工月度复核时可选的任务载荷
type IntegrationMonthlyFinalizePayload struct {
ArchiveMonth string `json:"archive_month"`
}
// IntegrationArchiveHandler 处理 Integration Log 每日归档与月度最终复核。
// IntegrationArchiveHandler 处理 Integration Log 每日归档
type IntegrationArchiveHandler struct {
service *auditarchive.Service
logger *zap.Logger
@@ -61,33 +56,6 @@ func (h *IntegrationArchiveHandler) HandleDaily(ctx context.Context, task *asynq
return nil
}
// HandleMonthlyFinalize 执行上一个完整自然月复核,或按任务载荷复核指定月份。
func (h *IntegrationArchiveHandler) HandleMonthlyFinalize(ctx context.Context, task *asynq.Task) error {
if h.service == nil {
return fmt.Errorf("Integration Log 归档服务未配置")
}
var err error
if len(task.Payload()) == 0 {
err = h.service.FinalizePreviousIntegrationMonth(ctx)
} else {
var payload IntegrationMonthlyFinalizePayload
if unmarshalErr := sonic.Unmarshal(task.Payload(), &payload); unmarshalErr != nil {
return fmt.Errorf("解析 Integration Log 月度复核任务载荷失败: %w", unmarshalErr)
}
month, parseErr := parseArchiveMonth(payload.ArchiveMonth)
if parseErr != nil {
return parseErr
}
err = h.service.FinalizeIntegrationMonth(ctx, month)
}
if err != nil {
h.logger.Error("Integration Log 月度最终版本复核失败,后续清理必须阻止", zap.Error(err))
return err
}
h.logger.Info("Integration Log 月度最终版本复核完成")
return nil
}
func parseArchiveDate(value string) (time.Time, error) {
location, err := time.LoadLocation(constants.AuditArchiveTimezone)
if err != nil {
@@ -99,15 +67,3 @@ func parseArchiveDate(value string) (time.Time, error) {
}
return date, nil
}
func parseArchiveMonth(value string) (time.Time, error) {
location, err := time.LoadLocation(constants.AuditArchiveTimezone)
if err != nil {
return time.Time{}, fmt.Errorf("加载 Integration Log 归档时区失败: %w", err)
}
month, err := time.ParseInLocation("2006-01", value, location)
if err != nil {
return time.Time{}, fmt.Errorf("解析 Integration Log 归档月份失败: %w", err)
}
return month, nil
}

View File

@@ -132,6 +132,29 @@ func (b *PollingBase) releaseCardTrafficSyncLock(_ context.Context, cardID uint,
}
}
// ensureMissingTask 根据数据库最新卡状态,仅在对应分片队列缺失时补入任务。
func (b *PollingBase) ensureMissingTask(ctx context.Context, cardID uint, taskType string) error {
card, err := b.iotCardStore.GetByID(ctx, cardID)
if err != nil {
return err
}
info, ok := b.configMgr.MergedTaskIntervals(card)[taskType]
if !ok || info.Interval <= 0 {
return nil
}
added, err := b.queueMgr.EnsureQueued(ctx, cardID, taskType, time.Now())
if err != nil {
return err
}
if added {
b.logger.Info("卡状态轮询补齐缺失套餐任务",
zap.Uint("card_id", cardID), zap.String("task_type", taskType))
}
return nil
}
// requeueCardAt 使用独立短超时上下文执行真正的 ZADD 重入队。
func (b *PollingBase) requeueCardAt(cardID uint, taskType string, nextCheckAt time.Time) error {
ctx, cancel := pollingFallbackContext()

View File

@@ -67,11 +67,15 @@ func (h *PollingCarddataHandler) Handle(ctx context.Context, task *asynq.Task) e
}
result, err := h.gateway.QueryFlow(ctx, &gateway.FlowQueryReq{CardNo: card.ICCID})
if err != nil {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 流量 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "查询流量失败", err)
}
if strings.TrimSpace(result.ICCID) == "" {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 流量 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "流量查询响应缺少 ICCID", nil)
}
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, true, attemptStartedAt); logErr != nil {

View File

@@ -66,11 +66,15 @@ func (h *PollingCardStatusHandler) Handle(ctx context.Context, task *asynq.Task)
}
result, err := h.gateway.QueryCardStatus(ctx, &gateway.CardStatusReq{CardNo: card.ICCID})
if err != nil {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 网络 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "查询卡状态失败", err)
}
if strings.TrimSpace(result.ICCID) == "" {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 网络 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "卡状态查询响应缺少 ICCID", nil)
}
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, true, attemptStartedAt); logErr != nil {
@@ -102,6 +106,12 @@ func (h *PollingCardStatusHandler) Handle(ctx context.Context, task *asynq.Task)
h.base.logger.Info("独立卡命中风险状态,已关闭轮询", zap.Uint("card_id", cardID), zap.String("gateway_extend", decision.GatewayExtend))
return nil
}
// 卡状态任务仍正常但套餐任务丢失时,仅补入缺失项;不改写已有套餐任务的执行时间。
h.base.invalidateCardCache(ctx, cardID)
if err := h.base.ensureMissingTask(ctx, cardID, constants.TaskTypePollingPackage); err != nil {
h.base.logger.Warn("卡状态轮询补齐套餐任务失败", zap.Uint("card_id", cardID), zap.Error(err))
}
return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus)
}

View File

@@ -38,11 +38,13 @@ func completeGatewayAttempt(ctx context.Context, repository *integrationlog.Repo
if repository == nil || attempt == nil {
return errors.New(errors.CodeInternalError, "Gateway Integration Log 尝试不存在")
}
completionCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer cancel()
result := constants.IntegrationResultFailed
if success {
result = constants.IntegrationResultSuccess
}
_, err := repository.Complete(ctx, attempt.IntegrationID, integrationlog.Completion{
_, err := repository.Complete(completionCtx, attempt.IntegrationID, integrationlog.Completion{
Result: result, DurationMS: time.Since(startedAt).Milliseconds(), StateChanged: false,
ResponseSummary: map[string]any{"success": success},
})

View File

@@ -58,11 +58,15 @@ func (h *PollingRealnameHandler) Handle(ctx context.Context, task *asynq.Task) e
}
result, err := h.gateway.QueryRealnameStatus(ctx, &gateway.CardStatusReq{CardNo: card.ICCID})
if err != nil {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 实名 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "查询实名状态失败", err)
}
if strings.TrimSpace(result.ICCID) == "" {
_ = completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt)
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, false, attemptStartedAt); logErr != nil {
return h.failAndRequeue(ctx, cardID, startedAt, "完成 Gateway 实名 Integration Log 失败", logErr)
}
return h.failAndRequeue(ctx, cardID, startedAt, "实名查询响应缺少 ICCID", nil)
}
if logErr := completeGatewayAttempt(ctx, h.integration, attempt, true, attemptStartedAt); logErr != nil {

View File

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

View File

@@ -0,0 +1,47 @@
## Context
代理在线充值现有 `OnlinePaymentPort` 抽象只有微信直连 H5/MWEB 与支付宝 WAP 两个 Adapter`payment-methods` 只调用 `wechat.Available``alipay.Available`。当前生效配置为富友时,微信 Adapter 因 `provider_type=fuiou` 判定不可用,导致只返回 `alipay`。富友本质是微信支付上游通道,`pkg/fuiou` 已具备 XML/GBK/双重 URL 编码/RSA 签名验签与回调解析能力,但缺少主扫下单与订单查询。参见 proposal.md - Why。
## Goals / Non-Goals
**Goals:**
- 对外支付方式枚举保持 `wechat` / `alipay`,不暴露 `fuiou`
- 富友配置下 `wechat` 走富友主扫统一下单,复用现有 `qr_content` 契约。
- 富友通道下主动查单通过 `/commonQuery` 收敛状态,不重建支付链接。
- 回调与确认校验接受 `channel=fuiou` 对应业务方式 `wechat`
**Non-Goals:**
- 不新增对外 `fuiou` 支付方式,不改前端枚举。
- 不接入富友退款、撤销、条码支付(商户扫用户)等其它交易类型。
- 不改支付宝通道(仍走直接支付宝 WAP
## Decisions
### 新增富友扫码 Adapter 而非扩展微信 Adapter
新增 `FuiouScanAdapter` 实现 `OnlinePaymentPort``Available` 判定 `provider_type==fuiou` 且富友字段完整;`CreatePaymentURL` 调主扫下单返回 `qr_code``Query` 调订单查询映射 `trans_stat`。备选方案是在 `WechatWebAdapter` 内部分支,但会混淆微信直连与富友的日志提供方、错误语义与恢复策略,故放弃。
### Adapter 选择改为配置感知
`OnlineCreationService``RecoverOnlinePaymentService``adapter()` 增加 `config` 参数:`wechat` 方法 + `provider_type==fuiou` 返回富友 Adapter否则返回微信 Adapter。恢复阶段富友与微信一致只查单不重建链接。
### 业务方式与渠道分离存储
`Payment.PaymentMethod` 与充值记录 `PaymentMethod` 保持 `wechat`(业务语义),充值记录 `PaymentChannel``fuiou`(实际渠道);`paymentMerchantIdentity` 在富友下返回 `FyMchntCd`。回调 `PaymentMethod=fuiou` 在确认入口归一化为业务方式 `wechat`,领域校验允许 `channel=fuiou` 映射到 `method=wechat`
### 复用现有富友回调
异步通知复用 `FuiouPayCallback``VerifyNotify``NotifyRequest` 字段与主扫通知报文一致,不做改动。
## Risks / Trade-offs
- [富友主扫 `mchnt_order_no` 必须全局唯一,重复会被拒绝] → 复用本地 `payment_no` 作为商户订单号,且恢复阶段只查单不重建链接。
- [富友查询 `trans_stat``9999`/空/`1010` 时状态未知] → 映射为 unknown保持待恢复继续查不确认收款也不关闭。
- [富友 `reserved_*` 字段不参与签名且渠道会新增] → 复用 `pkg/fuiou` 现有 `structToMap` 排除 reserved 前缀的签名规则。
- [回调无支付时间或金额不一致] → 现有确认用例已校验金额、配置身份与支付时间,富友金额用 `order_amt`(分)、时间用 `reserved_txn_fin_ts`
## Migration Plan
无数据库迁移、无新外部依赖。代码上线后,将生效支付配置切为富友即可使代理在线充值展示微信扫码;回滚为恢复生效配置为微信直连或回退代码,不改变既有数据语义。

View File

@@ -0,0 +1,31 @@
## Why
当前生效支付配置为富友(`provider_type=fuiou`)时,代理在线充值可用支付方式接口只返回 `alipay`,不返回 `wechat`。原因是现有微信 Adapter 只支持微信直连(`wechat`/`wechat_v2`),而富友虽然本质是微信支付上游通道,却未接入代理在线充值链路。本变更让富友主扫下单成为内部 `wechat` 通道,使代理在线充值在富友配置下也能展示并完成微信扫码支付。
## What Changes
- 代理在线充值 `payment-methods` 在富友配置完整时把 `wechat` 列为可用支付方式(对外枚举仍只有 `wechat``alipay`,不新增 `fuiou`)。
- `POST /api/admin/agent-recharges` 使用 `wechat` 创建时,若生效配置为富友,则调用富友主扫统一下单,返回 `qr_code` 作为支付链接。
- 新增富友主扫下单(`/preCreate`)与主动查单(`/commonQuery`)客户端能力;支付结果继续复用现有富友异步通知回调。
- 代理在线充值支付恢复(主动查单)在富友通道下通过 `commonQuery` 收敛状态,不重建链接。
- 支付确认校验允许富友渠道(`channel=fuiou`)对应业务支付方式 `wechat`
## Capabilities
### New Capabilities
(无)
### Modified Capabilities
- `agent-funds-commission`: 代理在线充值可用支付方式与创建行为在富友配置下按微信通道处理。
- `external-integration`: 富友主扫统一下单与订单查询的调用、状态映射与失败边界。
## Impact
- `internal/application/agentrecharge``OnlineCreationService``RecoverOnlinePaymentService`、确认用例渠道校验)
- `internal/domain/agentrecharge`(支付确认渠道一致性校验)
- `internal/infrastructure/payment`(新增富友扫码 Adapter
- `pkg/fuiou`(新增主扫下单与订单查询请求/响应)
- `internal/handler/callback`(富友回调支付方式归一化为业务方式 `wechat`
- 无数据库迁移、无新外部依赖;对外支付方式枚举不变。

View File

@@ -0,0 +1,29 @@
## ADDED Requirements
### Requirement: 代理在线充值可用支付方式按支付配置判定
系统 SHALL 按当前生效支付配置判定代理在线充值可用支付方式:微信直连(`wechat``wechat_v2`)配置完整,或富友(`fuiou`)配置完整时,返回 `wechat`;支付宝字段完整时返回 `alipay`。对外支付方式枚举 MUST 固定为 `wechat``alipay`MUST NOT 返回 `fuiou`。代理账号以 `wechat` 创建在线充值单时,若生效配置为富友,系统 MUST 使用富友主扫统一下单并将返回的二维码链接作为支付链接。
#### Scenario: 富友配置完整时微信可用
- **GIVEN** 当前生效支付配置 `provider_type=fuiou` 且富友机构号、商户号、终端号、私钥、公钥、API 地址、通知地址均非空
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `wechat` 的方式列表且不包含 `fuiou`
#### Scenario: 微信直连配置完整时微信可用
- **GIVEN** 当前生效支付配置为微信直连且对应字段完整
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `wechat` 的方式列表
#### Scenario: 支付宝字段完整时支付宝可用
- **GIVEN** 当前生效支付配置的支付宝应用 ID、应用私钥、支付宝公钥、通知地址均非空
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `alipay` 的方式列表
#### Scenario: 富友配置下微信创建走主扫下单
- **GIVEN** 当前生效支付配置为富友且字段完整
- **WHEN** 代理账号以 `wechat` 创建在线充值单
- **THEN** 系统调用富友主扫统一下单并返回二维码链接作为支付链接,本地充值单支付方式为 `wechat`、支付渠道为 `fuiou`

View File

@@ -0,0 +1,36 @@
## ADDED Requirements
### Requirement: 富友主扫统一下单与订单查询
系统 SHALL 通过富友主扫统一下单创建微信二维码支付并返回 `qr_code` 二维码链接;系统 SHALL 通过富友订单查询按 `trans_stat` 将状态映射为已支付、已关闭、待支付或未知,未知状态 MUST 保持待恢复。下单与查询失败 MUST 映射为项目稳定错误,已接入外部交互日志的调用保留脱敏结果。
#### Scenario: 主扫下单成功返回二维码链接
- **GIVEN** 富友支付配置完整
- **WHEN** 系统发起主扫统一下单且渠道返回成功
- **THEN** 系统返回 `qr_code` 作为支付链接,业务以该链接生成二维码
#### Scenario: 查询映射支付成功
- **WHEN** 富友订单查询返回 `trans_stat=SUCCESS`
- **THEN** 系统将状态映射为已支付并取得渠道交易号、金额与支付时间
#### Scenario: 查询映射已关闭
- **WHEN** 富友订单查询返回 `trans_stat``PAYERROR``CLOSED``REVOKED`
- **THEN** 系统将状态映射为已关闭
#### Scenario: 查询映射待支付
- **WHEN** 富友订单查询返回 `trans_stat``USERPAYING``NOTPAY`
- **THEN** 系统将状态映射为待支付
#### Scenario: 查询状态未知保持待恢复
- **WHEN** 富友订单查询返回系统错误、找不到交易或无法识别的 `trans_stat`
- **THEN** 系统将状态映射为未知并保持本地支付单待恢复,不得据此确认收款或关闭订单
#### Scenario: 下单失败返回稳定错误
- **WHEN** 富友主扫统一下单返回失败或请求结果未知
- **THEN** 系统返回项目稳定错误且已接入外部交互日志的调用记录脱敏结果

View File

@@ -0,0 +1,33 @@
## 1. 富友主扫下单与查询客户端
- [x] 1.1 在 `pkg/fuiou` 新增主扫统一下单请求/响应结构(`/preCreate`,含 `order_type``notify_url``reserved_expire_minute`,响应含 `qr_code`
- [x] 1.2 在 `pkg/fuiou` 新增主扫下单方法,复用 `Client.Sign``DoRequest`
- [x] 1.3 在 `pkg/fuiou` 新增订单查询请求/响应结构(`/commonQuery`,响应含 `trans_stat``order_amt``transaction_id``reserved_txn_fin_ts`
- [x] 1.4 在 `pkg/fuiou` 新增订单查询方法,复用 `Client.Sign``DoRequest`
- [x] 1.5 补充 `trans_stat` 到统一支付状态的映射SUCCESS→已支付PAYERROR/CLOSED/REVOKED→已关闭USERPAYING/NOTPAY→待支付其余→未知
## 2. 富友扫码 Adapter
- [x] 2.1 新增 `internal/infrastructure/payment/fuiou_scan.go``FuiouScanAdapter`,实现 `OnlinePaymentPort`
- [x] 2.2 `Available` 判定 `provider_type==fuiou` 且富友机构号/商户号/终端号/私钥/公钥/API 地址/通知地址完整
- [x] 2.3 `CreatePaymentURL` 调主扫下单并以 `qr_code` 作为 `QRContent`,接入外部交互日志
- [x] 2.4 `Query` 调订单查询并按映射返回统一查询结果
## 3. 代理在线充值 Adapter 选择与渠道事实
- [x] 3.1 `OnlineCreationService` 注入富友 Adapter`adapter()` 增加配置参数并按 `provider_type==fuiou` 分流
- [x] 3.2 `RecoverOnlinePaymentService` 同样注入并按配置分流,恢复阶段富友只查单不重建链接
- [x] 3.3 `paymentMerchantIdentity` 在富友下返回 `FyMchntCd`
- [x] 3.4 创建本地事实时 `PaymentChannel``fuiou``PaymentMethod` 仍存 `wechat`
## 4. 回调与确认校验
- [x] 4.1 `confirmAgentRechargePayment` 将回调 `PaymentMethod==fuiou` 归一化为业务方式 `wechat` 后再进入确认用例
- [x] 4.2 `domain.ValidatePaymentConfirmation` 允许 `channel=fuiou` 对应 `method=wechat`,保留其余一致性校验
## 5. 装配与验证
- [x] 5.1 `bootstrap/services.go` 注入富友扫码 Adapter 到在线创建与恢复服务
- [x] 5.2 `gofmt -w` 变更文件,`go build ./cmd/api ./cmd/worker` 通过
- [x] 5.3 `go run cmd/gendocs/main.go` 重新生成文档(如路由/DTO 有变化)
- [x] 5.4 `openspec validate --all` 通过

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非待审批、已关联、在线充值、原创建人/场景不可用及并发重复请求不创建第二实例。
- [x] 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-20

View File

@@ -0,0 +1,43 @@
## Context
见 proposal.md。退款申请创建已允许 0 元金额,但退款通过路径的公共金额校验拒绝 `<= 0`。企微终态消费者和受控人工审批均复用该校验。退款通过事务还会无条件调用钱包回款;代理主钱包回款命令不接受 0 元。
## Goals / Non-Goals
**Goals:**
- 让零金额退款在两条通过路径中保持一致的状态推进语义。
- 避免零金额钱包流水、余额变更或调用资金回款用例。
- 保持退款、订单、通知及既有后处理的事务与幂等边界。
**Non-Goals:**
- 不允许负数退款金额。
- 不修改退款申请、企微审批或订单的数据库结构。
- 不改变非零退款的校验、资金回款和后处理语义。
## Decisions
### 公共金额校验接受零、拒绝负数
将公共审批退款金额的下限从“必须大于零”调整为“不得小于零”,继续校验不超过申请金额和订单实收金额。这样企微终态与人工路径不会出现规则漂移。
备选方案是在企微消费者单独放行零金额;不采用,因为人工审批仍会拒绝,且重复了金额规则。
### 零金额在退款事务内跳过资金回款
当审批通过金额为零时,退款事务仍更新退款和订单状态、写入审计/通知/既有后处理事件,但跳过 `refundWalletPayment`。非零金额保留原调用。
备选方案是让钱包模块接受 0 元退款命令;不采用,因为会创建无资金含义的流水并放宽资金模块的金额不变量。
### 已失败的生产决策等待安全重试
现有 `approval:26:approved` 决策投递已经处于失败待重试,部署后由既有重试机制重新执行;若重试已耗尽,维护者按既有受控运维流程重放同一稳定事件,而不直接修改退款、订单或钱包数据。
## Risks / Trade-offs
- [0 元退款仍将订单标记为已退款] → 这是业务已确认语义;审计中保留审批退款金额为 0。
- [遗漏某条通过路径] → 两条路径复用公共校验,并在退款服务公共回款接缝处按金额分支。
- [旧失败任务已耗尽重试] → 发布后先核验决策投递状态;必要时以稳定事件 ID 受控重放。
## Migration Plan
1. 发布 API 与 Worker 二进制,不执行数据库迁移。
2. 维护者核验 `approval:26:approved` 的决策投递是否被重试并成功,以及退款单 `RF20260820170954264700` 是否更新。
3. 若未自动重试,由维护者受控重放同一审批终态事件并核验未产生 0 元钱包流水。
4. 回滚仅恢复旧二进制;已完成的零金额退款不回退业务事实。

View File

@@ -0,0 +1,24 @@
## Why
当前退款申请允许保存 0 元申请金额,但企微审批通过后的退款终态处理会拒绝该金额,导致通用审批已通过而退款申请持续待审批、反复重试。业务已确认 0 元退款是有效场景,需要将其作为不产生资金回款的退款完成。
## What Changes
- 允许 0 元退款申请在企微审批或受控人工审批通过后完成。
- 0 元退款通过时,将退款申请和关联订单推进为已通过/已退款,但不创建代理主钱包或资产钱包回款流水。
- 保持负数金额、超过申请金额或超过订单实收金额的审批退款金额无效。
## Capabilities
### New Capabilities
- `refund-approval`: 退款审批终态、订单状态和资金回款的业务语义。
### Modified Capabilities
- 无。
## Impact
- `internal/service/refund/` 的退款金额校验、企微审批终态处理和人工审批处理。
- 退款状态、订单支付状态、钱包流水和退款审批 Outbox 消费链路。
- 不新增 API、数据库迁移或第三方依赖。

View File

@@ -0,0 +1,30 @@
## Purpose
定义退款审批终态对退款、订单和资金回款事实的统一推进规则,确保零金额退款能够完成而不制造虚假的资金流水。
## ADDED Requirements
### Requirement: 零金额退款审批通过
系统 SHALL 接受金额为零的退款申请在企业微信审批或受控人工审批通过后完成;退款申请状态 SHALL 更新为已通过,关联订单支付状态 SHALL 更新为已退款,并记录审批通过时间和零金额审批退款金额。
#### Scenario: 企业微信通过零金额退款
- **WHEN** 本地关联的企业微信退款审批返回通过,且申请退款金额为零
- **THEN** 系统将退款申请和关联订单分别更新为已通过和已退款
#### Scenario: 人工通过零金额退款
- **WHEN** 启用的人工退款审批入口通过申请退款金额为零的退款申请
- **THEN** 系统将退款申请和关联订单分别更新为已通过和已退款
### Requirement: 零金额退款不产生资金回款
系统 SHALL 在零金额退款通过时不创建代理主钱包或资产钱包的退款回款流水,且不得调用资金回款处理。
#### Scenario: 钱包支付订单的零金额退款通过
- **WHEN** 钱包支付订单关联的零金额退款审批通过
- **THEN** 系统不写入任何该退款对应的钱包回款流水
### Requirement: 非法退款金额仍被拒绝
系统 MUST 拒绝负数审批退款金额、超过申请退款金额的审批退款金额,以及超过订单实收金额的审批退款金额。
#### Scenario: 负数退款金额
- **WHEN** 审批退款金额小于零
- **THEN** 系统拒绝完成退款且保持退款与订单原有状态

View File

@@ -0,0 +1,11 @@
## 1. 退款通过规则
- [x] 1.1 调整公共审批退款金额校验:接受零、拒绝负数,并保留申请金额和订单实收金额上限校验。
- [x] 1.2 在企微审批终态退款处理里,对零金额跳过钱包回款,同时保持退款、订单、审计、通知和后处理事件的既有事务语义。
- [x] 1.3 在受控人工退款审批路径复用相同的零金额回款跳过规则。
## 2. 验证与运维核验
- [x] 2.1 格式化变更的 Go 文件并执行 `go build ./cmd/api ./cmd/worker`
- [x] 2.2 执行 `openspec validate allow-zero-refund-approval --strict`,确认行为契约有效。
- [x] 2.3 为维护者记录生产发布后的核验项:重试或受控重放 `approval:26:approved`,确认退款单 `RF20260820170954264700` 与关联订单完成,且没有零金额钱包回款流水。

View File

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

View File

@@ -0,0 +1,68 @@
## Context
现有归档已按日生成 `tb_log_archive_run`,但 Integration Log 仅在月初终结整月 revision物理删除任务也只接受完整自然月。`tb_log_archive_run` 已按来源和归档日保存对象键、manifest、哈希、统计数和 `cleanup_started_at` / `cleaned_at`,可作为按日留存的断点账本。审计调查的在线边界当前通过每个来源的最大已清理范围计算;若出现清理日期空洞会错误隐藏仍在线的数据。
## Goals / Non-Goals
**Goals:**
- 每天将昨天及更早连续积压的归档日推进到可恢复、可校验、已物理清理的状态。
- 在任一日期归档、终结、对象校验或删除失败时保留在线数据,并使后续运行可安全续跑。
- 保持归档对象和 manifest 不删除,保持现有配置开关作为物理删除总闸门。
- 让在线查询边界只反映连续完成的日期。
**Non-Goals:**
- 不降低卡观测、轮询、Gateway 或 Integration Log 的写入频率。
- 不直接对生产表执行 SQL 删除,不绕过归档验证清理历史数据。
- 不改变归档 JSONL 格式、对象存储提供商或调查 API 路由。
## Decisions
### 以“归档日对”为最小留存单元
日留存以同一 Asia/Shanghai 日期的 Audit 与 Integration 两条归档账本为一个逻辑单元。任务处理上限固定为昨天,先确保两类归档存在;随后为 Integration 生成最终 revision确认该日不存在 `pending` 记录,并复用现有 manifest、对象 metadata、哈希和在线计数复核。
只有两个来源均通过验证后,才开始删除该日在线数据,顺序为 Audit Resource、Audit Event、Integration Log。删除仍按既有 1000 行批次进行;每个来源独立落 `cleanup_started_at``cleaned_at`,从而在进程中断后依据剩余计数续跑。
不采用“定时归档完成即直接 DELETE”的方案对象上传成功不等于可恢复且 Integration Log 的创建日可能仍有待完成记录。
### 按日期从早到晚推进,不跨越失败日
任务枚举昨天及更早仍未完成的归档日,从最早日期开始。某日失败、缺少归档或存在 pending Integration Log 时停止,不处理更晚日期。这样在线数据的已归档边界永远是连续前缀,历史积压在补齐归档后会由同一任务自动追赶。
不采用并行或跳过失败日期的方案:虽然可更快释放部分容量,但会产生在线数据日期空洞;当前调查响应只表达单一 `archived_before` 边界,不能正确描述空洞。
### 将 Integration 最终版由月度改为逐日形成
保留每日普通归档任务作为预归档;日留存任务在准备删除某日数据时对该日期执行最终归档。最终归档仍要求该日没有 pending Integration Log否则不删除并等待下次运行。原月度最终复核和月度留存调度移除避免与日留存争夺同一账本和产生不同的清理范围。
不采用固定延迟天数用户要求昨天及更早数据在可恢复后尽快清理pending 校验是每日期间的实际安全门禁。
### 查询边界按连续清理前缀计算
留存边界查询不再取任意来源的 `MAX(range_end)`。它必须从归档账本确认连续已清理日期,且多来源时间线使用两来源共同完成的连续边界;遇到未清理日期时停止。这样部分删除、失败重试或历史补档都不会把仍在线的旧日期标记为已归档。
### 使用独立文件记录留存结果与失败
新增独立的 Lumberjack/Zap 留存日志配置,默认路径为 Worker 工作目录下的 `logs/audit-retention.log`。日留存和日归档只向该日志写入日期推进、成功数量与耗时,或带日期、来源、失败分类的安全错误摘要;运营者无需在高噪声 `app.log` 中筛选。
不把诊断详情写入 `tb_log_archive_run.error_summary`。账本仍只保存归档状态、对象引用和清理断点;归档或留存失败的具体错误以独立文件为准。保留日志轮转和压缩,避免长期错误积压占满 Worker 磁盘。
不采用为每个日期新建单独文件的方案:一个按日轮转的专用流已经能按日期字段筛选,并避免文件数量随积压日期增长。
## Risks / Trade-offs
- [某日大量 pending 长期不终结,阻塞后续日期清理] → 记录日期、数量和原因;人工修复业务状态后由日任务续跑,不删除未完成数据。
- [单日数百万记录清理超过任务超时] → 维持小批次和来源级断点,任务重试时基于剩余记录续跑;按实测调整任务超时,不增大单事务范围。
- [归档或删除期间进程重启] → 账本的运行租约和清理标记保证归档可重做、删除可继续,且每次删除前重新验证剩余范围。
- [历史日期未建归档账本] → 日任务在最早缺失日期停止;先通过受控补归档形成完整连续账本,再允许自动物理清理。
- [日清理后数据库文件空间不立即下降] → 删除仅回收可复用空间;由维护者根据 PostgreSQL 运行状态决定后续 VACUUM 策略,不在 Worker 内执行高风险表重写。
- [独立留存日志写满磁盘或轮转配置错误] → 使用已有 Lumberjack 轮转、压缩和目录初始化机制;上线前核对 `logs/audit-retention.log` 的生成和轮转。
## Migration Plan
1. 发布包含日归档最终版、日留存和连续边界计算的 Worker/API 二进制,保持物理清理开关关闭。
2. 维护者核对历史归档账本的连续性按日期补齐缺失归档并以只读演练确认对象、manifest 和在线计数。
3. 维护者在低峰期启用 `JUNHONG_WORKER_AUDIT_RETENTION_CLEANUP_ENABLED=true` 并重启调度 Worker观察 `/opt/junhong_cmp/worker/logs/audit-retention.log`、账本断点和表大小。
4. 出现异常时关闭开关并重启调度 Worker已删除数据通过已验证的对象存储归档恢复未删除日期保留在线数据。

View File

@@ -0,0 +1,28 @@
## Why
`tb_audit_event``tb_audit_event_resource``tb_integration_log` 每日持续写入百万级记录;现有物理留存任务仅在月初清理上一个完整自然月,在线表、索引和备份在月内持续膨胀。生产需要在确认对象存储归档完整且可恢复后,按日清理在线日志数据。
## What Changes
- 将统一审计事件、审计资源快照和 Integration Log 的物理留存单位由完整自然月改为完整自然日。
- 每次日留存任务处理昨天及更早仍未完成清理的归档日;仅在对应 Audit 与 Integration 归档均成功、完整、可读取且满足最终版要求时,分批物理删除在线表中该日的数据。
- 将清理断点和查询留存边界改为按归档日表达,支持失败后从未完成日期续跑,且不删除未通过校验的日期。
- 保留对象存储归档和现有物理清理开关;不降低卡观测、轮询或外部交互日志的写入频率。
## Capabilities
### New Capabilities
无。
### Modified Capabilities
- `operations-audit`: 审计在线数据在逐日归档校验通过后按归档日物理留存,并向调查查询提供逐日清理边界。
- `external-integration`: 外部交互日志在逐日最终归档校验通过后按归档日物理留存。
## Impact
- Worker 定时任务、审计归档服务、月度留存处理器和留存边界查询。
- `tb_log_archive_run` 的清理断点语义,以及 `tb_audit_event``tb_audit_event_resource``tb_integration_log` 的物理删除范围。
- 对象存储归档读取、manifest 与哈希校验。
- 生产需更新 Worker 二进制、`JUNHONG_WORKER_AUDIT_RETENTION_CLEANUP_ENABLED` 及独立留存日志配置;日志写入 Worker 工作目录下的 `logs/audit-retention.log`,不新增外部 API不直接执行生产数据删除或迁移。

View File

@@ -0,0 +1,19 @@
## ADDED Requirements
### Requirement: 外部交互日志逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位物理留存 `tb_integration_log`。删除某日在线外部交互日志前,系统 MUST 为该日生成最终归档版本,确认不存在待完成记录,并验证归档对象、清单、日期范围、记录数量和校验摘要可作为恢复凭证。每次执行 SHALL 从最早尚未完成清理的归档日连续处理至昨天;任一日期不能形成或验证最终归档时,系统 MUST 不删除该日及更晚日期的在线外部交互日志。
#### Scenario: 最终归档通过后清理在线外部交互日志
- **GIVEN** 某已结束自然日的外部交互日志均已进入终态,且该日最终归档对象和清单校验成功
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日的在线外部交互日志并将该归档日标记为已清理,归档对象保持可用
#### Scenario: 存在待完成外部交互日志
- **GIVEN** 某归档日仍存在待完成的外部交互日志
- **WHEN** 日留存任务尝试处理该日期
- **THEN** 系统不删除该日及更晚日期的在线外部交互日志,并在独立日留存日志中记录该日期未形成最终归档的原因
#### Scenario: 历史积压按日期连续补清
- **GIVEN** 存在多个昨天及更早日期尚未完成物理清理
- **WHEN** 日留存任务执行且这些日期依次通过最终归档和完整性校验
- **THEN** 系统按日期从早到晚清理全部连续合格日期

View File

@@ -0,0 +1,41 @@
## ADDED Requirements
### Requirement: 审计在线数据逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位处理统一审计事件及其资源快照的在线留存。每次留存执行 SHALL 从最早尚未完成清理的归档日开始,连续处理至昨天;任一归档日未通过完整性校验时,系统 MUST 保留该日及其后续日期的在线审计数据,且不得将它们标记为已清理。系统 MUST 先删除该日的审计资源快照,再删除该日的审计事件,并保留对象存储归档对象及清单作为恢复凭证。
#### Scenario: 已验证日期完成物理清理
- **GIVEN** 某已结束自然日的审计归档成功,归档对象和清单可读取且与在线记录范围、数量和校验摘要一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日在线审计资源快照和审计事件,并将该归档日标记为已清理
#### Scenario: 归档校验失败阻断连续清理
- **GIVEN** 某尚未清理的归档日缺少归档、归档校验失败或清单与在线数据不一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统不删除该日或更晚日期的在线审计数据,并将失败信息写入独立日留存日志
#### Scenario: 失败后从未完成日期续跑
- **GIVEN** 日留存任务在某个归档日失败,较早日期已经标记为已清理
- **WHEN** 后续日留存任务再次执行
- **THEN** 系统从最早未完成清理的归档日继续处理,不重复删除已完成日期的数据
### Requirement: 日留存独立运行日志
系统 SHALL 将日留存的日期推进、清理成功和阻断失败写入独立轮转日志。该日志 MUST 位于 Worker 工作目录下配置的 `logs/audit-retention.log`,每条失败记录 MUST 包含归档日期、数据来源、失败分类和安全错误摘要。日留存和日归档路径 MUST NOT 将失败详情写入 `tb_log_archive_run.error_summary`
#### Scenario: 留存日期被阻断
- **GIVEN** 日留存处理某个归档日时发生缺失账本、pending 记录、归档校验或删除错误
- **WHEN** 系统停止该日期推进
- **THEN** `logs/audit-retention.log` 记录该日期、来源、失败分类和错误摘要
- **AND** 对应账本记录的 `error_summary` 不写入该失败详情
#### Scenario: 留存日期清理成功
- **GIVEN** 某归档日通过全部校验并完成在线数据物理删除
- **WHEN** 系统完成该日期处理
- **THEN** `logs/audit-retention.log` 记录归档日期、各来源清理数量和耗时
### Requirement: 审计调查留存边界连续
系统 SHALL 仅将连续完成物理清理的审计归档日期间公开为已归档边界。在线审计调查接口 MUST 继续允许查询任何尚未物理清理的较早日期,且不得因某个较晚日期已归档而将仍在线的日期错误标记为已归档。
#### Scenario: 清理链存在未完成日期
- **GIVEN** 某较早归档日尚未完成清理
- **WHEN** 后续日期的归档已成功生成
- **THEN** 审计调查接口仍将较早未清理日期视为在线可查询数据

View File

@@ -0,0 +1,19 @@
## 1. 按日归档与留存编排
- [x] 1.1 将 Integration Log 最终归档能力收敛为指定已结束自然日;保留 pending 记录时不产生最终 revision并通过独立留存日志输出原因而不更新 `error_summary`
- [x] 1.2 将现有月度留存服务改为逐日处理:枚举昨天及更早未完成日期,按日期从早到晚确保双来源归档、最终归档和完整性校验。
- [x] 1.3 复用现有分批删除和来源级清理断点,按 Audit Resource、Audit Event、Integration Log 的顺序物理删除单日在线数据;中断或失败时可从剩余记录续跑。
- [x] 1.4 遇到缺失账本、归档失败、对象或 manifest 校验失败、数量不一致或 pending Integration Log 时停止日期推进;不清理该日及后续日期,且不向 `tb_log_archive_run.error_summary` 写入失败详情。
## 2. Worker 调度与查询边界
- [x] 2.1 为 Worker 新增独立留存日志配置和目录初始化,默认输出 `logs/audit-retention.log`,复用现有轮转与压缩能力,并在退出时刷新日志。
- [x] 2.2 将 Worker 留存调度由月初任务改为每日任务,调整任务类型、唯一窗口、超时、处理器日志及归档注册,避免月度任务与日留存并发处理同一账本。
- [x] 2.3 保留 `JUNHONG_WORKER_AUDIT_RETENTION_CLEANUP_ENABLED` 作为物理删除总闸门;关闭时执行逐日只读校验,不写清理断点、不删除在线数据。
- [x] 2.4 改造审计和 Integration Log 在线查询留存边界,按连续已清理日期计算,并使混合时间线不将仍在线的空洞日期误报为已归档。
- [x] 2.5 更新留存模拟 CLI 或等价可执行验证入口以逐日范围验证归档、最终版、对象、manifest、数据库数量、清理断点及独立日志输出。
## 3. 运行说明与验证
- [x] 3.1 更新生产运行说明,明确日留存启用前的历史归档补齐、只读演练、开关启用、`/opt/junhong_cmp/worker/logs/audit-retention.log` 观察、账本断点、异常停用和归档恢复步骤。
- [x] 3.2 执行 `gofmt``go build ./cmd/api ./cmd/worker``go run cmd/gendocs/main.go``openspec validate daily-audit-log-retention --strict` 与日留存模拟验证;不新增或恢复自动化测试。(隔离库模拟待维护者在具备显式 `JUNHONG_*` 配置的环境执行确认)

View File

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

View File

@@ -0,0 +1,53 @@
## Context
见 proposal.md。现有实现以 Audit 和 Integration 归档对作为原子清理单元,且在生成 Integration 最终 revision 前要求整日不存在 `pending`。2026-08-18 的 458 条 Gateway 同步轮询遗留 pending 因而阻断约 314 万条 Integration Log 和后续日期的清理。
## Goals / Non-Goals
**Goals:**
- 只保留 pending Integration Log 在线,尽快释放已终态日志和审计数据。
- 保持删除前可恢复归档和逐批续跑语义。
- 防止旧 Gateway 轮询处理器静默制造新的永久 pending。
**Non-Goals:**
- 不直接修改历史 pending 的结果或删除它们。
- 不改变 Integration Log 的写入频率、JSONL 格式、对象存储格式或数据库 Schema。
- 不在 Worker 内执行 VACUUM 或表重写。
## Decisions
### 归档覆盖全日,删除只匹配终态
仍按创建日生成包含全部记录的归档 revision移除“存在 pending 即不得最终归档”的门禁。留存校验确认当前在线集合与该 revision 一致后,删除条件限定为 `result <> 'pending'`。这样每条已删除数据都在对象存储中存在恢复副本,而 pending 永不被删除。
不按结果分别创建归档文件:会改变既有归档格式和恢复语义,且没有必要。
### Audit 与 Integration 分来源推进
Audit 归档校验成功后独立完成 Audit 资源和事件删除Integration 的 pending 或残留不再使 Audit 退化为未清理。每个来源仍按自身失败日阻断该来源后续日期,避免跨越该来源的归档或校验故障。
不继续要求两来源同日原子清理:该设计正是容量积压的根因,且两类数据已有独立的归档账本和清理断点。
### 以剩余可删除记录决定 Integration 续跑
已经删除过终态记录的日期不因只剩 pending 而反复阻断日期扫描。若旧 pending 后续进入终态,则该日期重新生成 revision、验证当前在线集合并删除新增终态记录。查询仅检查可删除终态记录和对应清理断点避免对只含 pending 的大历史范围反复做全表工作。
### Gateway 失败分支不吞终结错误
三个旧轮询 Handler 对已创建的尝试使用受控收口:终结写入失败返回任务错误并保留重试信号;不再使用 `_ = completeGatewayAttempt(...)` 静默丢弃错误。任务取消时仍尝试使用不受取消影响的短生命周期上下文终结该条日志。
## Risks / Trade-offs
- [终态日志在归档和删除间新增或变化] → 归档后重新统计和校验;不一致时不删除并在下次重建 revision。
- [pending 永久不终结] → 仅保留这些记录,不能阻止终态数据和审计数据释放;通过调查接口和留存日志排查。
- [PostgreSQL 文件空间不立即缩小] → DELETE 回收空间供后续写入复用;维护者按运行手册评估 VACUUM不由 Worker 自动执行。
- [旧 pending 终结后遗漏清理] → 续跑选择条件识别清理后进入终态的记录,并重建该日 revision。
## Migration Plan
1. 发布 Worker 二进制,保持既有物理清理开关状态。
2. 若当前开关关闭,维护者在低峰期启用后重启调度 Worker。
3. 观察 `logs/audit-retention.log``tb_log_archive_run` 清理断点和数据库磁盘指标8 月 18 日的已终态记录应先被清理458 条 pending 保留。
4. 若归档或删除异常,关闭清理开关并重启 Worker已删除记录通过对应归档 revision 恢复,未删除记录仍在线。

View File

@@ -0,0 +1,26 @@
## Why
日留存将某日存在的任意 `pending` Integration Log 视为整日删除的阻断条件。Gateway 同步轮询的遗留 pending 记录使 2026-08-18 及所有后续日期无法清理,在线日志持续积压并带来数据库容量风险。
## What Changes
- 日留存继续归档全天 Integration Log但只将已进入公开终态的记录纳入可清理集合。
- `pending` 记录保留在线,且不得阻断同日终态记录和后续日期的归档、验证与清理。
- 已完成的 Gateway 轮询尝试在终结日志写入失败时不得静默忽略错误,避免新增永久 pending 记录。
## Capabilities
### New Capabilities
无。
### Modified Capabilities
- `external-integration`: 外部交互日志日留存从“整日无 pending 才能清理”改为“仅清理已终态记录并保留 pending”。
- `operations-audit`: 日留存不因 Integration Log 的 pending 记录阻断审计事件和资源快照的连续清理。
## Impact
- 代码:`internal/application/auditarchive/`、旧 Gateway 轮询任务处理器。
- Worker日留存会释放已终态的在线日志pending 继续可调查和补偿。
- 数据库:无迁移,不执行直接 SQL 数据修复。

View File

@@ -0,0 +1,59 @@
## MODIFIED Requirements
### Requirement: 外部交互日志逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位物理留存 `tb_integration_log`。删除一条已终态在线外部交互日志前,系统 MUST 已为其创建日生成可读取、可验证的归档对象和清单;该归档 MUST 包含该日当时全部外部交互日志,并验证日期范围、记录数量和校验摘要可作为恢复凭证。`pending` 记录 MUST 保留在线,不得被物理删除,且不得阻断同日已终态记录或更晚日期的归档、验证与物理清理。每次执行 SHALL 清理所有满足归档校验条件的已终态记录;某一来源的归档或校验失败时,系统 MUST 保留该来源当天及更晚日期的在线外部交互日志。
#### Scenario: 最终归档通过后清理在线外部交互日志
- **GIVEN** 某已结束自然日的外部交互归档对象和清单校验成功,且存在已终态在线外部交互日志
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日已终态的在线外部交互日志并保留归档对象
#### Scenario: 存在待完成外部交互日志
- **GIVEN** 某归档日仍存在待完成的外部交互日志
- **WHEN** 日留存任务尝试处理该日期
- **THEN** 系统保留该 pending 日志,但不得停止同日终态记录和更晚日期的清理
#### Scenario: 含 pending 的日期清理终态日志
- **GIVEN** 某已结束自然日同时存在已终态和 `pending` 的外部交互日志,且该日归档对象和清单校验成功
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日已终态的在线外部交互日志,保留 `pending` 记录,并保留归档对象
#### Scenario: 历史积压按日期连续补清
- **GIVEN** 存在多个昨天及更早日期尚未完成物理清理,且这些日期各自存在满足归档校验条件的已终态外部交互日志
- **WHEN** 日留存任务执行
- **THEN** 系统按日期从早到晚清理全部符合条件的终态日志
#### Scenario: pending 不阻断后续日期
- **GIVEN** 较早归档日只剩 `pending` 外部交互日志,且更晚归档日存在满足归档校验条件的已终态外部交互日志
- **WHEN** 日留存任务执行
- **THEN** 系统继续处理更晚归档日的已终态外部交互日志,不因较早日期的 pending 停止日期推进
#### Scenario: pending 后续终结
- **GIVEN** 某归档日的 pending 外部交互日志在该日首次留存后进入公开终态
- **WHEN** 后续日留存任务执行
- **THEN** 系统为包含该记录的当前在线集合重新验证归档,并删除该已终态记录
#### Scenario: 归档或校验失败保留同来源后续数据
- **GIVEN** 某归档日的外部交互归档缺失、归档校验失败或清单与在线数据不一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统不删除该日及更晚日期的在线外部交互日志,并将失败信息写入独立日留存日志
### Requirement: 外部失败边界
系统 SHALL 将第三方超时、渠道错误和无效响应转换为当前稳定的系统错误已接入外部交互日志的渠道同时保留脱敏结果。Gateway 同步轮询在已建立 pending 外部交互日志后无论查询成功、失败或任务上下文取消MUST 尝试将该尝试终结为公开终态;终结失败 MUST 作为任务失败返回或记录为可重试错误,不得静默忽略。
#### Scenario: 外部失败边界
- **GIVEN** 外部系统超时或返回失败
- **WHEN** 调用依赖该系统的操作
- **THEN** 客户端收到当前稳定错误;已接入外部交互日志的调用记录脱敏渠道结果
#### Scenario: Gateway 轮询查询失败
- **GIVEN** Gateway 同步轮询已建立 pending 外部交互日志且查询返回错误
- **WHEN** 轮询任务处理该错误
- **THEN** 系统将该日志终结为失败或结果未知的公开终态,并按既有策略重新入队
#### Scenario: Gateway 轮询终结日志失败
- **GIVEN** Gateway 同步轮询已建立 pending 外部交互日志且终结写入失败
- **WHEN** 轮询任务结束该次尝试
- **THEN** 系统记录并返回该终结失败,不得把该错误静默丢弃

View File

@@ -0,0 +1,25 @@
## MODIFIED Requirements
### Requirement: 审计在线数据逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位处理统一审计事件及其资源快照的在线留存。每次留存执行 SHALL 从最早尚未完成清理的归档日开始连续处理至昨天Integration Log 的 `pending` 记录不得阻断已通过完整性校验的审计事件及资源快照清理。任一审计归档日未通过完整性校验时,系统 MUST 保留该日及其后续日期的在线审计数据,且不得将它们标记为已清理。系统 MUST 先删除该日的审计资源快照,再删除该日的审计事件,并保留对象存储归档对象及清单作为恢复凭证。
#### Scenario: 已验证日期完成物理清理
- **GIVEN** 某已结束自然日的审计归档成功,归档对象和清单可读取且与在线记录范围、数量和校验摘要一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日在线审计资源快照和审计事件,并将该归档日标记为已清理
#### Scenario: Integration Log pending 不阻断审计清理
- **GIVEN** 某已结束自然日的审计归档已通过完整性校验,且同日存在 pending Integration Log
- **WHEN** 日留存任务处理该日期
- **THEN** 系统仍完成该日在线审计资源快照和审计事件的物理清理
#### Scenario: 归档校验失败阻断连续清理
- **GIVEN** 某尚未清理的审计归档日缺少归档、归档校验失败或清单与在线数据不一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统不删除该日或更晚日期的在线审计数据,并将失败信息写入独立日留存日志
#### Scenario: 失败后从未完成日期续跑
- **GIVEN** 日留存任务在某个审计归档日失败,较早日期已经标记为已清理
- **WHEN** 后续日留存任务再次执行
- **THEN** 系统从最早未完成清理的审计归档日继续处理,不重复删除已完成日期的数据

View File

@@ -0,0 +1,15 @@
## 1. 日留存按终态清理
- [x] 1.1 调整 Integration Log 归档与留存校验使归档保留全日快照、pending 不再阻断终态记录清理。
- [x] 1.2 将 Integration Log 删除条件限制为公开终态,并支持同一日期遗留 pending 后续终结时重建归档 revision 和续跑删除。
- [x] 1.3 让 Audit 留存独立于同日 Integration pending 完成归档校验和物理清理,并保持各来源的失败断点。
## 2. Gateway 轮询收口
- [x] 2.1 修改旧卡状态、流量和实名轮询 Handler使已创建尝试的终结失败不被忽略并在取消上下文下仍尝试终结。
## 3. 验证与运行说明
- [x] 3.1 格式化修改的 Go 文件并构建 `./cmd/worker`
- [ ] 3.2 以可控演练或最小验证确认:含 pending 的日期删除终态日志并继续后续日期pending 本身保留在线。
- [x] 3.3 更新生产运行说明,记录发布 Worker、启用既有清理开关、观察留存日志与 PostgreSQL 空间复用的步骤。

View File

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

View File

@@ -0,0 +1,32 @@
## Context
批量订购逐行调用后台订单创建,并在全部行处理结束后于同一事务保存任务结果和完成审计。当前完成审计以关联订单审计事件数量充当成功行数;该数量会包含失败记录和重试记录,可能大于输入总行数,违反审计表的批量统计约束并回滚任务完成状态。
## Goals / Non-Goals
**Goals:**
- 使任务完成审计使用本次处理结果的业务成功数和失败数。
- 保持任务结果、汇总计数和完成审计在同一事务提交。
**Non-Goals:**
- 不改变单笔订单、钱包余额不足或重复输入的业务规则。
- 不恢复或重放历史任务;历史数据由维护者按生产运维流程处置。
## Decisions
-`processRows` 返回的 `successCount``failCount` 作为任务完成审计的唯一批量统计来源。这些值与将写入任务记录的逐行结果来自同一次处理,满足审计表的计数不变量。
- 移除对关联订单审计事件数量的查询。审计事件数量会因失败记录和重试而膨胀,不能表示批量业务成功数。
## Risks / Trade-offs
- [历史待处理任务不自动恢复] → 保持现有状态,避免对已扣款订单重复执行;由维护者核对后受控结案。
- [完成审计仍依赖事务写入成功] → 使用与任务结果一致的统计值,消除本次约束失败原因,同时保留既有原子性。
## Migration Plan
1. 发布 API/Worker 中的修复后二进制,重点更新 Worker。
2. 在隔离环境验证部分成功任务返回完成状态与一致计数。
3. 生产历史任务由维护者在核对订单和源文件后手工结案;不重投原 Asynq 消息。
4. 若发布异常,回滚 Worker 二进制;该修复不含迁移。

View File

@@ -0,0 +1,23 @@
## Why
资产套餐批量订购在部分行因钱包余额不足或输入重复而失败时,完成审计将审计事件条数误作成功行数,可能违反审计批量统计约束并回滚任务终态。已实际创建的订单因此在任务查询中显示为待处理且无结果,既不能反映业务结果,也有误重试风险。
## What Changes
- 批量订购任务完成时,使用本次逐行处理得到的成功数和失败数写入完成审计。
- 确保部分成功、部分失败的批量订购任务能持久化为已完成,并可查询真实汇总和逐行结果。
## Capabilities
### New Capabilities
无。
### Modified Capabilities
- `package-lifecycle`: 异步批量订购任务必须持久化并返回逐行处理后的终态统计,即使部分资产因余额不足或输入校验失败而未能购包。
## Impact
- `internal/task/asset_package_batch_order.go` 中的批量订购完成与审计统计。
- 不新增 API、数据表或依赖既有批量订购详情将正确展示部分完成结果。

View File

@@ -0,0 +1,17 @@
## MODIFIED Requirements
### Requirement: 批量操作可追踪
系统 SHALL 为同步批量分配和调价直接返回处理结果;对异步批量订购返回任务标识并提供状态查询。异步批量订购完成后,系统 MUST 持久化每个输入行的成功或失败结果及与其一致的总数、成功数和失败数;部分资产因余额不足、资产校验或重复输入失败不得阻止任务进入完成终态。
#### Scenario: 批量操作可追踪
- **GIVEN** 操作者提交非空且有权处理的资源集合
- **WHEN** 创建批量操作
- **THEN** 同步操作直接返回结果;异步订购返回任务标识且可查询处理状态
#### Scenario: 批量订购部分失败后查询结果
- **GIVEN** 异步批量订购中的部分资产已成功创建订单,其他资产因钱包余额不足或输入重复失败
- **WHEN** Worker 完成全部输入行的处理
- **THEN** 任务状态为已完成,逐行结果保留成功订单与失败原因,且总数等于成功数与失败数之和

View File

@@ -0,0 +1,8 @@
## 1. 批量订购完成统计修复
- [x] 1.1 修改批量订购完成审计,使其使用逐行处理产生的成功数和失败数,不再统计关联审计事件数量。
- [x] 1.2 格式化修改的 Go 文件并构建 Worker验证部分成功任务的完成审计统计满足任务总行数约束。
## 2. 规格验证
- [x] 2.1 运行 OpenSpec 校验,确认变更工件一致。

View File

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

View File

@@ -0,0 +1,66 @@
## Context
见 proposal.md。当前到期处理在事务提交后投递异步套餐激活任务同时以 goroutine 发起停机检查;两条路径没有共同顺序或共享状态结论。生产中已观察到新套餐激活后 1.9 至 15.7 秒仍写入 `no_package` 停机。
现有 `ActivationService` 已负责待生效主套餐激活及激活后的复机回调;`StopResumeService` 已集中停复机条件。保留这两个职责边界,不增加新的状态表或外部依赖。
## Goals / Non-Goals
**Goals:**
- 同一载体到期接续期间不以过时权益快照停机。
- 没有后续可生效套餐时保留现有停机语义。
- 激活投递或执行暂时失败时宁可交由既有恢复/轮询兜底,不误判为无套餐。
- 给维护者提供基于当前事实的只读候选筛选与受控补偿流程。
**Non-Goals:**
- 不新增自动批量复机、管理接口或数据库迁移。
- 不把 Redis/Asynq 运行排障改造成业务状态机。
- 不改变手动停机、运营商风险停机或实名策略语义。
## Decisions
### 1. 将停机决策后移到套餐激活任务的最终重新评估
到期处理找到后续待生效主套餐时,只投递激活任务,不并发启动停机检查。激活任务完成后以该载体的最新套餐状态调用现有停机评估:成功激活时不触发无套餐停机;队首套餐因实名等业务条件不能激活时,才按当前事实评估停机。
这样复用现有激活锁、套餐条件和停复机服务,不让两个异步流程各自读取一次状态。
备选方案:在到期处理内同步完成套餐激活。未采用,因为到期扫描单轮最多处理大量套餐,同步串行激活会放大调度 Worker 的数据库和 Redis 锁持有时间。
### 2. 激活投递失败不是无套餐结论
仅在明确不存在后续主套餐时,允许到期处理立即发起停机评估。存在后续套餐但投递失败时记录错误并保留当前网络状态,由既有孤儿套餐恢复扫描和套餐轮询重试;不得回退到并发停机。
备选方案:投递失败立即停机。未采用,因为它重新引入“接续结果未知被当作无套餐”的错误路径。
### 3. 存量恢复与代码修复分离
代码发布后,维护者用只读查询筛选“有效主套餐、未耗尽、实名满足、`no_package` 停机”的候选卡;每张或受控小批次先查询运营商实时状态,再调用既有复机能力。补偿不直接更新数据库,也不自动对全部候选执行。
备选方案:发布时自动批量复机全部候选。未采用,因为本地状态不足以代替运营商风险、销户和人工业务判断。
### 4. 轮询只作兜底并单独核验运行覆盖
停机卡配置的 `polling:package` 仍是后备恢复机制,不能作为避免竞态的主路径。维护者分别核验调度 Worker 心跳、Redis 分片队列深度、Asynq 队列积压与消费者吞吐;这些运行指标只用于定位未补偿原因。
### 5. 卡状态轮询仅补齐缺失的套餐任务
生产日志确认存在卡状态和流量任务仍持续执行、但同一卡的 `polling:package` 队列项缺失的情形。卡状态轮询成功且未命中风险停机后,按数据库最新卡状态重新匹配轮询配置;若该卡应参与套餐检查,则以 Redis `ZADD NX` 仅在对应分片 Sorted Set 不存在该卡时立即补入套餐任务。已存在任务不得覆盖其 score也不新增重复任务。
该做法不将卡状态轮询变为停复机决策路径:它只修复缺失的套餐调度项,最终仍由既有 `PollingPackageHandler → EvaluateAndAct` 判断权益、实名、停机原因和运营商结果。
## Risks / Trade-offs
- [激活任务投递失败时卡暂不因套餐到期自动停机] → 保留孤儿恢复扫描和套餐轮询;记录失败并在运行排障中核验任务恢复。
- [激活任务内增加最终停机评估] → 保持幂等条件和现有 Redis 锁;同一载体状态变化必须重读数据库。
- [存量候选包含运营商侧不可复机卡] → 补偿前逐卡或受控小批查询运营商状态,拒绝风险停机、销户或其他非本地套餐原因。
- [轮询吞吐不足延迟兜底] → 发布前后分别记录队列深度、消费速率和最近检查时间分布,不以单次数据库字段判断 Worker 健康。
- [卡状态轮询补齐套餐任务] → 仅使用 `ZADD NX` 补入缺失成员,不改写已有任务的执行时间;最终复机仍受既有权益、实名、停机原因和运营商校验约束。
## Migration Plan
1. 在隔离环境验证“旧套餐到期 + 后续套餐可激活”不会产生停机调用,及“无后续/不可激活”仍进入停机评估。
2. 按生产运行说明发布 Worker不需要数据库迁移。
3. 维护者观察一轮套餐到期处理、激活任务与轮询队列,确认没有新的“有效套餐 + `no_package` 停机”记录。
4. 使用只读筛选得到存量候选,按受控批次核验运营商状态后调用既有复机;失败项保留失败原因并人工处理。
5. 若发布后出现意外停机语义,回滚 Worker 二进制;存量补偿不得直接写库。

View File

@@ -0,0 +1,28 @@
## Why
生产诊断发现套餐接续期间会并发投递下一套餐激活任务和无套餐停机检查。停机检查可在新套餐生效后才完成并回写停机,造成卡已具备有效套餐、实名和未耗尽流量却仍以 `no_package` 停机;当前已检出 117 张此类卡,且 13 张在新套餐生效后 120 秒内被停机。
## What Changes
- 将主套餐到期后的“接续下一套餐”和“无有效套餐停机”收敛为同一条按载体串行的决策流程,禁止待生效套餐接续期间以旧快照发起停机。
- 只有在确认不存在可生效的后续主套餐,或后续套餐经业务校验确定不能生效后,才执行无套餐停机检查。
- 套餐接续完成后重新依据当前卡与套餐事实判断停复机,确保有效套餐卡不遗留 `no_package` 停机状态。
- 为既有异常卡提供只读筛选条件和受控人工补偿步骤;补偿前必须重新核验套餐、流量、实名和运营商实时状态,不引入自动批量复机。
- 排查并记录轮询调度、入队与消费覆盖情况;当卡状态轮询仍正常但 `polling:package` 队列项缺失时,仅补入缺失的套餐任务,确保套餐轮询仍可作为即时接续失败后的兜底,不将 Redis/Asynq 运行状态误判为业务终态。
## Capabilities
### New Capabilities
- 无。
### Modified Capabilities
- `package-lifecycle`: 主套餐到期、后续套餐接续及无套餐停机必须基于同一载体的最新权益事实,避免接续成功后错误停机。
## Impact
- 受影响代码:`internal/polling/package_activation_handler.go``internal/polling/queue_manager.go``internal/service/package/activation_service.go``internal/service/iot_card/stop_resume_service.go` 及相关轮询任务装配。
- 受影响运行单元:调度 Worker、Asynq 套餐激活任务、套餐轮询任务。
- 不新增 HTTP API、数据库 Schema 或第三方依赖。
- 生产补偿仅由维护者执行受控外部复机Agent 仅提供只读诊断与候选清单。

View File

@@ -0,0 +1,35 @@
## MODIFIED Requirements
### Requirement: 套餐状态流转
系统 SHALL 按当前套餐和套餐使用状态控制上架、订购、激活、失效与到期处理。主套餐到期时,系统 MUST 先确定同一载体是否存在待生效的后续主套餐:存在时,后续套餐激活与停复机重新评估 MUST 由同一条顺序流程完成;系统 MUST NOT 依据后续套餐激活前的无套餐快照发起停机。后续套餐成功生效后,系统 MUST 依据最新套餐、流量和实名事实重新判断卡网络状态,且不得遗留 `no_package` 停机。不存在后续套餐或后续套餐经业务校验不能生效时,系统 SHALL 按现有停机规则评估卡状态。后续套餐激活结果未知或任务投递失败不得被当作无后续套餐处理并据此停机,系统 SHALL 保留既有激活恢复与轮询兜底路径。
#### Scenario: 套餐状态流转
- **GIVEN** 套餐或使用记录处于允许的前置状态
- **WHEN** 执行状态操作
- **THEN** 仅发生一次允许的状态变化;不满足前置状态时返回业务错误
#### Scenario: 到期主套餐接续后续套餐
- **GIVEN** 某载体的当前主套餐到期,且存在满足激活条件的待生效后续主套餐
- **WHEN** 系统处理该主套餐到期
- **THEN** 系统先完成后续套餐激活并按最新权益事实重新评估停复机,且不得因到期前的无套餐快照对该载体发起 `no_package` 停机
#### Scenario: 到期主套餐无后续可生效套餐
- **GIVEN** 某载体的当前主套餐到期,且不存在后续主套餐或队首后续主套餐不满足激活条件
- **WHEN** 系统完成该套餐到期处理
- **THEN** 系统按当前套餐、流量和实名事实执行既有停机评估
#### Scenario: 后续套餐激活结果未知
- **GIVEN** 某载体的当前主套餐到期,存在待生效后续主套餐,但激活任务投递或执行结果暂时未知
- **WHEN** 系统处理该套餐到期
- **THEN** 系统不得将该未知结果视为不存在后续套餐而依据旧快照发起停机,并保留既有激活恢复与套餐轮询兜底
#### Scenario: 卡状态轮询发现缺失的套餐任务
- **GIVEN** 启用轮询的卡匹配套餐检查配置,且其 `polling:package` 分片队列项因异常缺失
- **WHEN** 卡状态轮询成功完成且未命中风险停机
- **THEN** 系统基于最新卡状态仅补入缺失的套餐任务,不改写已存在套餐任务的执行时间;后续套餐任务仍按既有停复机条件评估该卡

View File

@@ -0,0 +1,21 @@
## 1. 到期接续顺序
- [x] 1.1 调整主套餐到期处理:明确区分“无后续主套餐”“已提交后续激活”和“后续激活投递失败”,仅无后续主套餐可立即进入停机评估。
- [x] 1.2 调整套餐激活任务:在后续套餐激活完成或确认暂不能激活后,重读载体当前权益并调用既有停机评估;移除与激活任务并发的停机 goroutine。
- [x] 1.3 保持激活后的既有复机回调、Redis 锁、审计与 Integration Log 语义,不新增状态表、接口或自动批量复机。
## 2. 运行兜底与可观测性
- [x] 2.1 核对停机卡匹配的 `polling:package` 配置及生命周期回调入队路径,修正本次改动涉及的入队或重排遗漏,确保其仅作为接续异常后的兜底。
- [x] 2.2 为激活任务的最终停机评估保留可按载体、套餐使用记录和结果关联的中文结构化日志;不得记录敏感外部载荷。
## 3. 验证
- [ ] 3.1 在隔离环境人工演练“旧套餐到期且后续套餐可激活”“无后续套餐”“后续套餐暂不能激活”“激活投递失败”四种路径,核对停复机调用与套餐最终状态。
- [x] 3.2 执行 `gofmt -w <changed-go-files>``go build ./cmd/api ./cmd/worker``openspec validate fix-package-expiry-stop-race --strict`
## 4. 生产受控恢复
- [x] 4.1 维护者通过只读查询重新筛选“有效主套餐、流量未耗尽、实名满足且 `no_package` 停机”的候选卡,并记录筛选时点与数量。
- [ ] 4.2 维护者按受控小批次查询运营商实时状态,仅对可复机卡调用既有复机能力;不得直接更新卡网络状态或停机原因。
- [ ] 4.3 维护者核验调度 Worker 心跳、Redis 分片队列、Asynq 队列和消费者吞吐,并在发布后观察是否仍产生“有效套餐 + `no_package` 停机”记录。

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,66 @@
- **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: 代理在线充值可用支付方式按支付配置判定
系统 SHALL 按当前生效支付配置判定代理在线充值可用支付方式:微信直连(`wechat``wechat_v2`)配置完整,或富友(`fuiou`)配置完整时,返回 `wechat`;支付宝字段完整时返回 `alipay`。对外支付方式枚举 MUST 固定为 `wechat``alipay`MUST NOT 返回 `fuiou`。代理账号以 `wechat` 创建在线充值单时,若生效配置为富友,系统 MUST 使用富友主扫统一下单并将返回的二维码链接作为支付链接。
#### Scenario: 富友配置完整时微信可用
- **GIVEN** 当前生效支付配置 `provider_type=fuiou` 且富友机构号、商户号、终端号、私钥、公钥、API 地址、通知地址均非空
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `wechat` 的方式列表且不包含 `fuiou`
#### Scenario: 微信直连配置完整时微信可用
- **GIVEN** 当前生效支付配置为微信直连且对应字段完整
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `wechat` 的方式列表
#### Scenario: 支付宝字段完整时支付宝可用
- **GIVEN** 当前生效支付配置的支付宝应用 ID、应用私钥、支付宝公钥、通知地址均非空
- **WHEN** 代理账号查询可用支付方式
- **THEN** 系统返回包含 `alipay` 的方式列表
#### Scenario: 富友配置下微信创建走主扫下单
- **GIVEN** 当前生效支付配置为富友且字段完整
- **WHEN** 代理账号以 `wechat` 创建在线充值单
- **THEN** 系统调用富友主扫统一下单并返回二维码链接作为支付链接,本地充值单支付方式为 `wechat`、支付渠道为 `fuiou`
## 可达操作索引
本节只用于入口导航,不是行为 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

@@ -28,7 +28,7 @@
### Requirement: 外部失败边界
系统 SHALL 将第三方超时、渠道错误和无效响应转换为当前稳定的系统错误;已接入外部交互日志的渠道同时保留脱敏结果。
系统 SHALL 将第三方超时、渠道错误和无效响应转换为当前稳定的系统错误;已接入外部交互日志的渠道同时保留脱敏结果。Gateway 同步轮询在已建立 pending 外部交互日志后无论查询成功、失败或任务上下文取消MUST 尝试将该尝试终结为公开终态;终结失败 MUST 作为任务失败返回或记录为可重试错误,不得静默忽略。
#### Scenario: 外部失败边界
@@ -36,6 +36,18 @@
- **WHEN** 调用依赖该系统的操作
- **THEN** 客户端收到当前稳定错误;已接入外部交互日志的调用记录脱敏渠道结果
#### Scenario: Gateway 轮询查询失败
- **GIVEN** Gateway 同步轮询已建立 pending 外部交互日志且查询返回错误
- **WHEN** 轮询任务处理该错误
- **THEN** 系统将该日志终结为失败或结果未知的公开终态,并按既有策略重新入队
#### Scenario: Gateway 轮询终结日志失败
- **GIVEN** Gateway 同步轮询已建立 pending 外部交互日志且终结写入失败
- **WHEN** 轮询任务结束该次尝试
- **THEN** 系统记录并返回该终结失败,不得把该错误静默丢弃
### Requirement: 外部调用重试边界
系统 SHALL 仅对 Gateway 客户端超时、连接失败和 DNS 失败自动重试默认最多重试两次且每次重新签名HTTP 非 200、响应解析失败、Gateway 业务失败和调用方取消不重试。企业微信审批只有确认尚未调用提交接口的失败可释放后重试,提交结果未知时不得盲目重建审批。
@@ -66,6 +78,80 @@
- **WHEN** 渠道再次发送相同业务事实
- **THEN** 系统返回渠道可接受响应且不重复推进卡状态
### Requirement: 富友主扫统一下单与订单查询
系统 SHALL 通过富友主扫统一下单创建微信二维码支付并返回 `qr_code` 二维码链接;系统 SHALL 通过富友订单查询按 `trans_stat` 将状态映射为已支付、已关闭、待支付或未知,未知状态 MUST 保持待恢复。下单与查询失败 MUST 映射为项目稳定错误,已接入外部交互日志的调用保留脱敏结果。
#### Scenario: 主扫下单成功返回二维码链接
- **GIVEN** 富友支付配置完整
- **WHEN** 系统发起主扫统一下单且渠道返回成功
- **THEN** 系统返回 `qr_code` 作为支付链接,业务以该链接生成二维码
#### Scenario: 查询映射支付成功
- **WHEN** 富友订单查询返回 `trans_stat=SUCCESS`
- **THEN** 系统将状态映射为已支付并取得渠道交易号、金额与支付时间
#### Scenario: 查询映射已关闭
- **WHEN** 富友订单查询返回 `trans_stat``PAYERROR``CLOSED``REVOKED`
- **THEN** 系统将状态映射为已关闭
#### Scenario: 查询映射待支付
- **WHEN** 富友订单查询返回 `trans_stat``USERPAYING``NOTPAY`
- **THEN** 系统将状态映射为待支付
#### Scenario: 查询状态未知保持待恢复
- **WHEN** 富友订单查询返回系统错误、找不到交易或无法识别的 `trans_stat`
- **THEN** 系统将状态映射为未知并保持本地支付单待恢复,不得据此确认收款或关闭订单
#### Scenario: 下单失败返回稳定错误
- **WHEN** 富友主扫统一下单返回失败或请求结果未知
- **THEN** 系统返回项目稳定错误且已接入外部交互日志的调用记录脱敏结果
### Requirement: 外部交互日志逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位物理留存 `tb_integration_log`。删除一条已终态在线外部交互日志前,系统 MUST 已为其创建日生成可读取、可验证的归档对象和清单;该归档 MUST 包含该日当时全部外部交互日志,并验证日期范围、记录数量和校验摘要可作为恢复凭证。`pending` 记录 MUST 保留在线,不得被物理删除,且不得阻断同日已终态记录或更晚日期的归档、验证与物理清理。每次执行 SHALL 清理所有满足归档校验条件的已终态记录;某一来源的归档或校验失败时,系统 MUST 保留该来源当天及更晚日期的在线外部交互日志。
#### Scenario: 最终归档通过后清理在线外部交互日志
- **GIVEN** 某已结束自然日的外部交互归档对象和清单校验成功,且存在已终态在线外部交互日志
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日已终态的在线外部交互日志并保留归档对象
#### Scenario: 存在待完成外部交互日志
- **GIVEN** 某归档日仍存在待完成的外部交互日志
- **WHEN** 日留存任务尝试处理该日期
- **THEN** 系统保留该 pending 日志,但不得停止同日终态记录和更晚日期的清理
#### Scenario: 含 pending 的日期清理终态日志
- **GIVEN** 某已结束自然日同时存在已终态和 `pending` 的外部交互日志,且该日归档对象和清单校验成功
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日已终态的在线外部交互日志,保留 `pending` 记录,并保留归档对象
#### Scenario: 历史积压按日期连续补清
- **GIVEN** 存在多个昨天及更早日期尚未完成物理清理,且这些日期各自存在满足归档校验条件的已终态外部交互日志
- **WHEN** 日留存任务执行
- **THEN** 系统按日期从早到晚清理全部符合条件的终态日志
#### Scenario: pending 不阻断后续日期
- **GIVEN** 较早归档日只剩 `pending` 外部交互日志,且更晚归档日存在满足归档校验条件的已终态外部交互日志
- **WHEN** 日留存任务执行
- **THEN** 系统继续处理更晚归档日的已终态外部交互日志,不因较早日期的 pending 停止日期推进
#### Scenario: pending 后续终结
- **GIVEN** 某归档日的 pending 外部交互日志在该日首次留存后进入公开终态
- **WHEN** 后续日留存任务执行
- **THEN** 系统为包含该记录的当前在线集合重新验证归档,并删除该已终态记录
#### Scenario: 归档或校验失败保留同来源后续数据
- **GIVEN** 某归档日的外部交互归档缺失、归档校验失败或清单与在线数据不一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统不删除该日及更晚日期的在线外部交互日志,并将失败信息写入独立日留存日志
## 可达操作索引
本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。

View File

@@ -46,6 +46,54 @@
- **WHEN** 该操作的审计事件构造、校验或持久化失败
- **THEN** 系统提交或返回该业务操作原本的结果,并以请求关联标识、动作编码和资源标识记录审计失败
### Requirement: 审计在线数据逐日物理留存
系统 SHALL 以 Asia/Shanghai 已结束的自然日为单位处理统一审计事件及其资源快照的在线留存。每次留存执行 SHALL 从最早尚未完成清理的归档日开始连续处理至昨天Integration Log 的 `pending` 记录不得阻断已通过完整性校验的审计事件及资源快照清理。任一审计归档日未通过完整性校验时,系统 MUST 保留该日及其后续日期的在线审计数据,且不得将它们标记为已清理。系统 MUST 先删除该日的审计资源快照,再删除该日的审计事件,并保留对象存储归档对象及清单作为恢复凭证。
#### Scenario: 已验证日期完成物理清理
- **GIVEN** 某已结束自然日的审计归档成功,归档对象和清单可读取且与在线记录范围、数量和校验摘要一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统分批删除该日在线审计资源快照和审计事件,并将该归档日标记为已清理
#### Scenario: Integration Log pending 不阻断审计清理
- **GIVEN** 某已结束自然日的审计归档已通过完整性校验,且同日存在 pending Integration Log
- **WHEN** 日留存任务处理该日期
- **THEN** 系统仍完成该日在线审计资源快照和审计事件的物理清理
#### Scenario: 归档校验失败阻断连续清理
- **GIVEN** 某尚未清理的归档日缺少归档、归档校验失败或清单与在线数据不一致
- **WHEN** 日留存任务处理该日期
- **THEN** 系统不删除该日或更晚日期的在线审计数据,并将失败信息写入独立日留存日志
#### Scenario: 失败后从未完成日期续跑
- **GIVEN** 日留存任务在某个归档日失败,较早日期已经标记为已清理
- **WHEN** 后续日留存任务再次执行
- **THEN** 系统从最早未完成清理的归档日继续处理,不重复删除已完成日期的数据
### Requirement: 日留存独立运行日志
系统 SHALL 将日留存的日期推进、清理成功和阻断失败写入独立轮转日志。该日志 MUST 位于 Worker 工作目录下配置的 `logs/audit-retention.log`,每条失败记录 MUST 包含归档日期、数据来源、失败分类和安全错误摘要。日留存和日归档路径 MUST NOT 将失败详情写入 `tb_log_archive_run.error_summary`
#### Scenario: 留存日期被阻断
- **GIVEN** 日留存处理某个归档日时发生缺失账本、pending 记录、归档校验或删除错误
- **WHEN** 系统停止该日期推进
- **THEN** `logs/audit-retention.log` 记录该日期、来源、失败分类和错误摘要
- **AND** 对应账本记录的 `error_summary` 不写入该失败详情
#### Scenario: 留存日期清理成功
- **GIVEN** 某归档日通过全部校验并完成在线数据物理删除
- **WHEN** 系统完成该日期处理
- **THEN** `logs/audit-retention.log` 记录归档日期、各来源清理数量和耗时
### Requirement: 审计调查留存边界连续
系统 SHALL 仅将连续完成物理清理的审计归档日期间公开为已归档边界。在线审计调查接口 MUST 继续允许查询任何尚未物理清理的较早日期,且不得因某个较晚日期已归档而将仍在线的日期错误标记为已归档。
#### Scenario: 清理链存在未完成日期
- **GIVEN** 某较早归档日尚未完成清理
- **WHEN** 后续日期的归档已成功生成
- **THEN** 审计调查接口仍将较早未清理日期视为在线可查询数据
## 可达操作索引
本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。

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`(退回退款申请)。
### 换货管理

View File

@@ -8,7 +8,7 @@
### Requirement: 套餐状态流转
系统 SHALL 按当前套餐和套餐使用状态控制上架、订购、激活、失效与到期处理。
系统 SHALL 按当前套餐和套餐使用状态控制上架、订购、激活、失效与到期处理。主套餐到期时,系统 MUST 先确定同一载体是否存在待生效的后续主套餐:存在时,后续套餐激活与停复机重新评估 MUST 由同一条顺序流程完成;系统 MUST NOT 依据后续套餐激活前的无套餐快照发起停机。后续套餐成功生效后,系统 MUST 依据最新套餐、流量和实名事实重新判断卡网络状态,且不得遗留 `no_package` 停机。不存在后续套餐或后续套餐经业务校验不能生效时,系统 SHALL 按现有停机规则评估卡状态。后续套餐激活结果未知或任务投递失败不得被当作无后续套餐处理并据此停机,系统 SHALL 保留既有激活恢复与轮询兜底路径。
#### Scenario: 套餐状态流转
@@ -16,9 +16,33 @@
- **WHEN** 执行状态操作
- **THEN** 仅发生一次允许的状态变化;不满足前置状态时返回业务错误
#### Scenario: 到期主套餐接续后续套餐
- **GIVEN** 某载体的当前主套餐到期,且存在满足激活条件的待生效后续主套餐
- **WHEN** 系统处理该主套餐到期
- **THEN** 系统先完成后续套餐激活并按最新权益事实重新评估停复机,且不得因到期前的无套餐快照对该载体发起 `no_package` 停机
#### Scenario: 到期主套餐无后续可生效套餐
- **GIVEN** 某载体的当前主套餐到期,且不存在后续主套餐或队首后续主套餐不满足激活条件
- **WHEN** 系统完成该套餐到期处理
- **THEN** 系统按当前套餐、流量和实名事实执行既有停机评估
#### Scenario: 后续套餐激活结果未知
- **GIVEN** 某载体的当前主套餐到期,存在待生效后续主套餐,但激活任务投递或执行结果暂时未知
- **WHEN** 系统处理该套餐到期
- **THEN** 系统不得将该未知结果视为不存在后续套餐而依据旧快照发起停机,并保留既有激活恢复与套餐轮询兜底
#### Scenario: 卡状态轮询发现缺失的套餐任务
- **GIVEN** 启用轮询的卡匹配套餐检查配置,且其 `polling:package` 分片队列项因异常缺失
- **WHEN** 卡状态轮询成功完成且未命中风险停机
- **THEN** 系统基于最新卡状态仅补入缺失的套餐任务,不改写已存在套餐任务的执行时间;后续套餐任务仍按既有停复机条件评估该卡
### Requirement: 批量操作可追踪
系统 SHALL 为同步批量分配和调价直接返回处理结果;对异步批量订购返回任务标识并提供状态查询。
系统 SHALL 为同步批量分配和调价直接返回处理结果;对异步批量订购返回任务标识并提供状态查询。异步批量订购完成后,系统 MUST 持久化每个输入行的成功或失败结果及与其一致的总数、成功数和失败数;部分资产因余额不足、资产校验或重复输入失败不得阻止任务进入完成终态。
#### Scenario: 批量操作可追踪
@@ -26,6 +50,12 @@
- **WHEN** 创建批量操作
- **THEN** 同步操作直接返回结果;异步订购返回任务标识且可查询处理状态
#### Scenario: 批量订购部分失败后查询结果
- **GIVEN** 异步批量订购中的部分资产已成功创建订单,其他资产因钱包余额不足或输入重复失败
- **WHEN** Worker 完成全部输入行的处理
- **THEN** 任务状态为已完成,逐行结果保留成功订单与失败原因,且总数等于成功数与失败数之和
### Requirement: 授权页面禁止重复选择套餐
系统 SHALL 使代理系列授权页面能够区分目标店铺已授权和未授权套餐;已授权套餐 MUST 以不可新增的状态返回,首次创建系列授权和既有系列新增套餐均适用。

View File

@@ -0,0 +1,32 @@
# refund-approval Specification
## Purpose
定义退款审批终态对退款、订单和资金回款事实的统一推进规则,确保零金额退款能够完成而不制造虚假的资金流水。
## Requirements
### Requirement: 零金额退款审批通过
系统 SHALL 接受金额为零的退款申请在企业微信审批或受控人工审批通过后完成;退款申请状态 SHALL 更新为已通过,关联订单支付状态 SHALL 更新为已退款,并记录审批通过时间和零金额审批退款金额。
#### Scenario: 企业微信通过零金额退款
- **WHEN** 本地关联的企业微信退款审批返回通过,且申请退款金额为零
- **THEN** 系统将退款申请和关联订单分别更新为已通过和已退款
#### Scenario: 人工通过零金额退款
- **WHEN** 启用的人工退款审批入口通过申请退款金额为零的退款申请
- **THEN** 系统将退款申请和关联订单分别更新为已通过和已退款
### Requirement: 零金额退款不产生资金回款
系统 SHALL 在零金额退款通过时不创建代理主钱包或资产钱包的退款回款流水,且不得调用资金回款处理。
#### Scenario: 钱包支付订单的零金额退款通过
- **WHEN** 钱包支付订单关联的零金额退款审批通过
- **THEN** 系统不写入任何该退款对应的钱包回款流水
### Requirement: 非法退款金额仍被拒绝
系统 MUST 拒绝负数审批退款金额、超过申请退款金额的审批退款金额,以及超过订单实收金额的审批退款金额。
#### Scenario: 负数退款金额
- **WHEN** 审批退款金额小于零
- **THEN** 系统拒绝完成退款且保持退款与订单原有状态

View File

@@ -13,6 +13,7 @@ type DirectoryResult struct {
TempDir string
AppLogDir string
AccessLogDir string
RetentionLogDir string
Fallbacks []string
}
@@ -27,6 +28,7 @@ func EnsureDirectories(cfg *config.Config, logger *zap.Logger) (*DirectoryResult
{cfg.Storage.TempDir, "storage.temp_dir", &result.TempDir},
{filepath.Dir(cfg.Logging.AppLog.Filename), "logging.app_log.filename", &result.AppLogDir},
{filepath.Dir(cfg.Logging.AccessLog.Filename), "logging.access_log.filename", &result.AccessLogDir},
{filepath.Dir(cfg.Logging.RetentionLog.Filename), "logging.retention_log.filename", &result.RetentionLogDir},
}
for _, dir := range directories {

View File

@@ -83,6 +83,7 @@ type LoggingConfig struct {
Development bool `mapstructure:"development"` // 启用开发模式(美化输出)
AppLog LogRotationConfig `mapstructure:"app_log"` // 应用日志配置
AccessLog LogRotationConfig `mapstructure:"access_log"` // HTTP 访问日志配置
RetentionLog LogRotationConfig `mapstructure:"retention_log"` // Worker 日留存日志配置
}
// LogRotationConfig Lumberjack 日志轮转配置
@@ -168,7 +169,7 @@ type GatewayConfig struct {
type WorkerConfig struct {
Role string `mapstructure:"role"` // Worker 运行角色all、leader、consumer
InstanceName string `mapstructure:"instance_name"` // Worker 实例名称,用于多实例日志区分
AuditRetentionCleanupEnabled bool `mapstructure:"audit_retention_cleanup_enabled"` // 是否启用审计日志月度物理清理
AuditRetentionCleanupEnabled bool `mapstructure:"audit_retention_cleanup_enabled"` // 是否启用审计日志物理清理
}
// ApprovalConfig 审批新旧入口切换配置。
@@ -291,12 +292,18 @@ func (c *Config) Validate() error {
if c.Logging.AccessLog.Filename == "" {
return fmt.Errorf("invalid configuration: logging.access_log.filename: must be non-empty valid file path")
}
if c.Logging.RetentionLog.Filename == "" {
return fmt.Errorf("invalid configuration: logging.retention_log.filename: must be non-empty valid file path")
}
if c.Logging.AppLog.MaxSize < 1 || c.Logging.AppLog.MaxSize > 1000 {
return fmt.Errorf("invalid configuration: logging.app_log.max_size: size out of range (current value: %d, expected: 1-1000 MB)", c.Logging.AppLog.MaxSize)
}
if c.Logging.AccessLog.MaxSize < 1 || c.Logging.AccessLog.MaxSize > 1000 {
return fmt.Errorf("invalid configuration: logging.access_log.max_size: size out of range (current value: %d, expected: 1-1000 MB)", c.Logging.AccessLog.MaxSize)
}
if c.Logging.RetentionLog.MaxSize < 1 || c.Logging.RetentionLog.MaxSize > 1000 {
return fmt.Errorf("invalid configuration: logging.retention_log.max_size: size out of range (current value: %d, expected: 1-1000 MB)", c.Logging.RetentionLog.MaxSize)
}
// 中间件验证
if c.Middleware.RateLimiter.Max <= 0 {

View File

@@ -65,6 +65,12 @@ logging:
max_backups: 3
max_age: 7
compress: true
retention_log:
filename: "logs/audit-retention.log"
max_size: 100
max_backups: 3
max_age: 7
compress: true
# 任务队列配置
queue:
@@ -137,7 +143,7 @@ polling_auto_trigger:
worker:
role: "all"
instance_name: ""
# 完整自然月灰度验收通过前必须保持关闭
# 日留存只读演练验收通过前必须保持关闭
audit_retention_cleanup_enabled: false
# 审批新旧入口切换配置

View File

@@ -94,6 +94,11 @@ func bindEnvVariables(v *viper.Viper) {
"logging.access_log.max_backups",
"logging.access_log.max_age",
"logging.access_log.compress",
"logging.retention_log.filename",
"logging.retention_log.max_size",
"logging.retention_log.max_backups",
"logging.retention_log.max_age",
"logging.retention_log.compress",
"queue.concurrency",
"queue.retry_max",
"queue.timeout",

View File

@@ -5,10 +5,8 @@ const (
TaskTypeAuditDailyArchive = "audit:daily:archive"
// TaskTypeIntegrationDailyArchive 表示 Integration Log 每日冷归档任务。
TaskTypeIntegrationDailyArchive = "integration:daily:archive"
// TaskTypeIntegrationMonthlyFinalize 表示 Integration Log 月度最终版本复核任务。
TaskTypeIntegrationMonthlyFinalize = "integration:monthly:finalize"
// TaskTypeAuditMonthlyRetention 表示审计日志月度物理清理任务。
TaskTypeAuditMonthlyRetention = "audit:monthly:retention"
// TaskTypeAuditDailyRetention 表示审计日志日留存任务。
TaskTypeAuditDailyRetention = "audit:daily:retention"
// AuditArchiveSource 表示 Audit Event 与 Event Resource 归档数据源。
AuditArchiveSource = "audit"
@@ -32,6 +30,6 @@ const (
// ArchiveStatusFailed 表示归档任务执行失败并等待重试。
ArchiveStatusFailed = "failed"
// AuditRetentionDeleteBatchSize 表示月度物理清理单批删除上限。
// AuditRetentionDeleteBatchSize 表示物理清理单批删除上限。
AuditRetentionDeleteBatchSize = 1000
)

View File

@@ -298,7 +298,7 @@ func QueueForTaskType(taskType string) string {
return QueueDataCleanup
case TaskTypeDailyTrafficFlush:
return QueueDailyTrafficFlush
case TaskTypeAuditDailyArchive, TaskTypeIntegrationDailyArchive, TaskTypeIntegrationMonthlyFinalize, TaskTypeAuditMonthlyRetention:
case TaskTypeAuditDailyArchive, TaskTypeIntegrationDailyArchive, TaskTypeAuditDailyRetention:
return QueueDataCleanup
case TaskTypeOutboxDeliver:
return QueueOutboxDeliver

92
pkg/fuiou/scan.go Normal file
View File

@@ -0,0 +1,92 @@
package fuiou
import (
"fmt"
"time"
"go.uber.org/zap"
)
// PreCreate 主扫统一下单C扫B生成用户扫码支付的二维码
// orderType: 订单类型WECHAT=微信主扫、ALIPAY=支付宝主扫。
// reservedExpireMinute: 订单有效期分钟reserved 字段,不参与签名)。
func (c *Client) PreCreate(orderNo, amount, goodsDesc, termIP, orderType, reservedExpireMinute string) (*PreCreateResponse, error) {
req := &PreCreateRequest{
Version: "1",
InsCd: c.InsCd,
MchntCd: c.MchntCd,
TermId: c.TermId,
RandomStr: generateRandomStr(),
OrderType: orderType,
MchntOrderNo: orderNo,
OrderAmt: amount,
GoodsDesc: goodsDesc,
TermIp: termIP,
TxnBeginTs: time.Now().Format("20060102150405"),
NotifyUrl: c.NotifyURL,
ReservedExpireMinute: reservedExpireMinute,
}
sign, err := c.Sign(req)
if err != nil {
return nil, fmt.Errorf("签名失败: %w", err)
}
req.Sign = sign
var resp PreCreateResponse
if err := c.DoRequest("/preCreate", req, &resp); err != nil {
return nil, fmt.Errorf("请求富友失败: %w", err)
}
if resp.ResultCode != "000000" {
c.logger.Error("富友主扫下单失败",
zap.String("order_no", orderNo),
zap.String("result_code", resp.ResultCode),
zap.String("result_msg", resp.ResultMsg),
)
return nil, fmt.Errorf("富友主扫下单失败: %s", resp.ResultMsg)
}
c.logger.Info("富友主扫下单成功",
zap.String("order_no", orderNo),
zap.String("order_type", orderType),
)
return &resp, nil
}
// CommonQuery 主动查询订单状态。
// orderType: 订单类型,须与下单时一致。
func (c *Client) CommonQuery(orderNo, orderType string) (*CommonQueryResponse, error) {
req := &CommonQueryRequest{
Version: "1",
InsCd: c.InsCd,
MchntCd: c.MchntCd,
TermId: c.TermId,
RandomStr: generateRandomStr(),
OrderType: orderType,
MchntOrderNo: orderNo,
}
sign, err := c.Sign(req)
if err != nil {
return nil, fmt.Errorf("签名失败: %w", err)
}
req.Sign = sign
var resp CommonQueryResponse
if err := c.DoRequest("/commonQuery", req, &resp); err != nil {
return nil, fmt.Errorf("请求富友失败: %w", err)
}
if resp.ResultCode != "000000" {
c.logger.Error("富友订单查询失败",
zap.String("order_no", orderNo),
zap.String("result_code", resp.ResultCode),
zap.String("result_msg", resp.ResultMsg),
)
return nil, fmt.Errorf("富友订单查询失败: %s", resp.ResultMsg)
}
return &resp, nil
}

View File

@@ -4,6 +4,12 @@ package fuiou
import "encoding/xml"
// 主扫统一下单与订单查询的订单类型常量
const (
OrderTypeWechat = "WECHAT" // 微信主扫
OrderTypeAlipay = "ALIPAY" // 支付宝主扫
)
// WxPreCreateRequest wxPreCreate 下单请求3.3接口)
// 所有非 reserved 字段必须出现在 XML 中(含空值),否则富友验签不通过
type WxPreCreateRequest struct {
@@ -77,3 +83,69 @@ type NotifyResponse struct {
ResultCode string `xml:"result_code"` // 结果码
ResultMsg string `xml:"result_msg"` // 结果消息
}
// PreCreateRequest 主扫统一下单请求C扫B
// 所有非 reserved 字段必须出现在 XML 中(含空值),否则富友验签不通过
type PreCreateRequest struct {
XMLName xml.Name `xml:"xml"`
Version string `xml:"version"` // 版本号: 1.0
InsCd string `xml:"ins_cd"` // 机构号
MchntCd string `xml:"mchnt_cd"` // 商户号
TermId string `xml:"term_id"` // 终端号
RandomStr string `xml:"random_str"` // 随机字符串
Sign string `xml:"sign"` // 签名
OrderType string `xml:"order_type"` // 订单类型: WECHAT=微信主扫
MchntOrderNo string `xml:"mchnt_order_no"` // 商户订单号
CurrType string `xml:"curr_type"` // 货币类型(可选)
OrderAmt string `xml:"order_amt"` // 订单金额(分)
TermIp string `xml:"term_ip"` // 终端IP
TxnBeginTs string `xml:"txn_begin_ts"` // 交易起始时间,格式 yyyyMMddHHmmss
NotifyUrl string `xml:"notify_url"` // 回调地址
GoodsDesc string `xml:"goods_des"` // 商品描述
GoodsDetail string `xml:"goods_detail"` // 商品详情(可选)
GoodsTag string `xml:"goods_tag"` // 商品标记(可选,非 reserved参与签名
AddnInf string `xml:"addn_inf"` // 附加数据(可选)
ReservedExpireMinute string `xml:"reserved_expire_minute"` // 订单有效期分钟reserved不参与签名
}
// PreCreateResponse 主扫统一下单响应
type PreCreateResponse struct {
ResultCode string `xml:"result_code"` // 结果码: 000000=成功
ResultMsg string `xml:"result_msg"` // 结果消息
InsCd string `xml:"ins_cd"` // 机构号
MchntCd string `xml:"mchnt_cd"` // 商户号
RandomStr string `xml:"random_str"` // 随机字符串
Sign string `xml:"sign"` // 签名
QrCode string `xml:"qr_code"` // 二维码内容
ReservedFyTraceNo string `xml:"reserved_fy_trace_no"` // 富友流水号
}
// CommonQueryRequest 订单查询请求
type CommonQueryRequest struct {
XMLName xml.Name `xml:"xml"`
Version string `xml:"version"` // 版本号: 1.0
InsCd string `xml:"ins_cd"` // 机构号
MchntCd string `xml:"mchnt_cd"` // 商户号
TermId string `xml:"term_id"` // 终端号
RandomStr string `xml:"random_str"` // 随机字符串
Sign string `xml:"sign"` // 签名
OrderType string `xml:"order_type"` // 订单类型
MchntOrderNo string `xml:"mchnt_order_no"` // 商户订单号
}
// CommonQueryResponse 订单查询响应
type CommonQueryResponse struct {
ResultCode string `xml:"result_code"` // 结果码: 000000=成功
ResultMsg string `xml:"result_msg"` // 结果消息
InsCd string `xml:"ins_cd"` // 机构号
MchntCd string `xml:"mchnt_cd"` // 商户号
RandomStr string `xml:"random_str"` // 随机字符串
Sign string `xml:"sign"` // 签名
OrderType string `xml:"order_type"` // 订单类型
MchntOrderNo string `xml:"mchnt_order_no"` // 商户订单号
TransStat string `xml:"trans_stat"` // 交易状态
OrderAmt string `xml:"order_amt"` // 订单金额(分)
TransactionId string `xml:"transaction_id"` // 交易流水号
ReservedTxnFinTs string `xml:"reserved_txn_fin_ts"` // 交易完成时间reserved
ReservedFyTraceNo string `xml:"reserved_fy_trace_no"` // 富友流水号reserved
}

17
pkg/logger/retention.go Normal file
View File

@@ -0,0 +1,17 @@
package logger
import (
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
)
// NewRetentionLogger 创建仅供 Worker 日留存使用的独立轮转日志。
func NewRetentionLogger(level string, config LogRotationConfig) (*zap.Logger, func() error) {
encoder := zapcore.NewJSONEncoder(zapcore.EncoderConfig{
TimeKey: "time", LevelKey: "level", MessageKey: "msg", CallerKey: "caller",
EncodeTime: zapcore.ISO8601TimeEncoder, EncodeLevel: zapcore.CapitalLevelEncoder,
EncodeCaller: zapcore.ShortCallerEncoder,
})
logger := zap.New(zapcore.NewCore(encoder, zapcore.AddSync(newLumberjackLogger(config)), parseLevel(level)), zap.AddCaller())
return logger, logger.Sync
}

Some files were not shown because too many files have changed in this diff Show More