diff --git a/config/config.go b/config/config.go index 0079b29..85f6f96 100644 --- a/config/config.go +++ b/config/config.go @@ -14,6 +14,7 @@ type Config struct { AlipaySandbox AlipaySandboxConfig `mapstructure:"alipay_sandbox"` Wechat WechatConfig `mapstructure:"wechat"` QuerySync QuerySyncConfig `mapstructure:"query_sync"` + Reconcile ReconcileConfig `mapstructure:"reconcile"` // Biz 通用:任意业务系统(jiu / dudu / …)在 config 的 biz. 下声明即可接入,无需改代码。 Biz map[string]BizSystemConfig `mapstructure:"biz"` // Accounts 多账户配置注册表(v2):凭证不写死配置,只存 env 前缀,真值运行时从环境变量取。 @@ -98,6 +99,23 @@ type QuerySyncConfig struct { MaxAgeMin int `mapstructure:"max_age_min"` // 只查创建时间在该分钟数内的待支付单 } +// ReconcileConfig v2 后台守护/对账周期任务开关与间隔。 +type ReconcileConfig struct { + Enabled bool `mapstructure:"enabled"` + OrderTTLMin int `mapstructure:"order_ttl_min"` // pending 订单存活 TTL(分钟),超则关闭 + ExpireEverySec int `mapstructure:"expire_every_sec"` // 过期清理间隔 + SyncEverySec int `mapstructure:"sync_every_sec"` // 查单对账间隔 + UsageEverySec int `mapstructure:"usage_every_sec"` // 用量快照刷新间隔 + SpotCheckEverySec int `mapstructure:"spot_check_every_sec"` // 已付抽查间隔 + SpotCheckWindowMin int `mapstructure:"spot_check_window_min"` // 抽查回溯窗(分钟) + OrphanEverySec int `mapstructure:"orphan_every_sec"` // crypto 孤儿扫描间隔(Task 6) + OrphanWindowMin int `mapstructure:"orphan_window_min"` // 孤儿扫描回溯窗(Task 6) + // —— 以下两项为 P4 T3 opus review 追加义务(退款修复扫描 + 卡滞退款告警), + // 本期归到本任务(对账主体)一起装配,brief 原表未列。 + RefundApplyEverySec int `mapstructure:"refund_apply_every_sec"` // 退款修复扫描 + 卡滞告警 共用间隔 + RefundStuckWarnMin int `mapstructure:"refund_stuck_warn_min"` // 退款卡滞 processing/manual_pending 告警阈值(分钟) +} + var C Config func Load() { @@ -130,6 +148,17 @@ func Load() { viper.SetDefault("query_sync.enabled", true) viper.SetDefault("query_sync.interval_sec", 30) viper.SetDefault("query_sync.max_age_min", 30) + viper.SetDefault("reconcile.enabled", true) + viper.SetDefault("reconcile.order_ttl_min", 60) + viper.SetDefault("reconcile.expire_every_sec", 300) + viper.SetDefault("reconcile.sync_every_sec", 30) + viper.SetDefault("reconcile.usage_every_sec", 60) + viper.SetDefault("reconcile.spot_check_every_sec", 300) + viper.SetDefault("reconcile.spot_check_window_min", 180) + viper.SetDefault("reconcile.orphan_every_sec", 300) + viper.SetDefault("reconcile.orphan_window_min", 180) + viper.SetDefault("reconcile.refund_apply_every_sec", 300) + viper.SetDefault("reconcile.refund_stuck_warn_min", 30) if err := viper.ReadInConfig(); err != nil { log.Println("[config] 未找到 config.yaml,使用默认值 + 环境变量") diff --git a/internal/provider/crypto/crypto.go b/internal/provider/crypto/crypto.go index 33b7520..0576132 100644 --- a/internal/provider/crypto/crypto.go +++ b/internal/provider/crypto/crypto.go @@ -73,6 +73,10 @@ func WithHTTPClient(c *http.Client) Option { return func(p *Provider func WithReservationLoader(l ReservationLoader) Option { return func(p *Provider) { p.loader = l } } func WithNow(f func() time.Time) Option { return func(p *Provider) { p.now = f } } +// SetReservationLoader 构造后注入冷启动预留源(装配期 main 在 BuildRegistry 之后调用: +// loader 依赖 OrderStore,而注册表构造不便传 store)。非并发安全,仅启动期单线程调用。 +func (p *Provider) SetReservationLoader(l ReservationLoader) { p.loader = l } + func New(accts *accounts.Registry, opts ...Option) *Provider { p := &Provider{ accts: accts, diff --git a/internal/reconcile/refund_apply.go b/internal/reconcile/refund_apply.go new file mode 100644 index 0000000..fbfb561 --- /dev/null +++ b/internal/reconcile/refund_apply.go @@ -0,0 +1,119 @@ +package reconcile + +import ( + "context" + "log" + "time" + + "github.com/wangjia/pay/internal/model" + "github.com/wangjia/pay/internal/store" +) + +// RefundApplyTask 是「RefundApply 修复扫描」:P4 T3 opus review 追加的义务(P4 退款 +// 代码本体在主 checkout,未落到本 worktree;此处只对本 worktree 已有的 +// store.RefundStore/OrderStore 接口(P4 T2,在 base 里)建自愈扫描,设计为能在合并 +// P4 后继续工作)。 +// +// 动机:退款成功的崩溃窗口 —— RefundStore.MarkRefundStatus 把某笔退款翻成 +// succeeded 后,调用方在再调 OrderStore.ApplyRefundToOrder 前进程崩溃/网络抖动, +// 订单状态卡在 paid(或旧的 partially_refunded),与「已实际退款成功」的事实脱节。 +// 本任务周期重算每个候选订单的 succeeded 退款之和,按既有 ApplyRefundToOrder 的 +// 条件 UPDATE 语义重新 apply 一次 —— 状态已一致时 UPDATE 影响 0 行,天然幂等。 +// +// 候选订单 = distinct(有 succeeded 退款的订单) ∪ (当前处于 refunding/partially_refunded +// 态的订单):前者直接命中"退款成功但订单未跟上"的崩溃窗口;后者兜住"已在退款流程 +// 中、但后续又有退款 succeeded 未被重算"的情形。 +func RefundApplyTask(orders *store.OrderStore, refunds *store.RefundStore, limit int) func(ctx context.Context) error { + return func(ctx context.Context) error { + succeededNos, err := refunds.ListDistinctOutTradeNosByStatus(model.RefundSucceeded, limit) + if err != nil { + return err + } + refundingOrders, err := orders.ListOrdersByStatus( + []model.OrderStatusV2{model.OrderRefundingV2, model.OrderPartRefundedV2}, limit) + if err != nil { + return err + } + + seen := make(map[string]bool, len(succeededNos)+len(refundingOrders)) + candidates := make([]string, 0, len(succeededNos)+len(refundingOrders)) + for _, no := range succeededNos { + if !seen[no] { + seen[no] = true + candidates = append(candidates, no) + } + } + for i := range refundingOrders { + no := refundingOrders[i].OutTradeNo + if !seen[no] { + seen[no] = true + candidates = append(candidates, no) + } + } + + for _, no := range candidates { + if err := reapplyRefundState(orders, refunds, no); err != nil { + log.Printf("[reconcile] 退款修复扫描 out_trade_no=%s: %v", no, err) + } + } + return nil + } +} + +// reapplyRefundState 对单个订单重算 succeeded 退款之和并按需重新 apply 状态转移。 +// 只有目标态与当前态不同才真正调用 ApplyRefundToOrder(避免每轮扫描都打"翻转"日志噪声); +// 没有 succeeded 退款(sum==0)的订单跳过 —— 它不属于本扫描要修的窗口。 +func reapplyRefundState(orders *store.OrderStore, refunds *store.RefundStore, outTradeNo string) error { + succ, err := refunds.RefundSum(outTradeNo, model.RefundSucceeded) + if err != nil { + return err + } + if succ <= 0 { + return nil + } + o, err := orders.GetOrder(outTradeNo) + if err != nil { + return err + } + + fully := succ >= o.AmountMinor + next := model.OrderPartRefundedV2 + if fully { + next = model.OrderRefundedV2 + } + if o.Status == next { + return nil // 已一致,无需自愈 + } + + flipped, err := orders.ApplyRefundToOrder(outTradeNo, fully) + if err != nil { + return err + } + if flipped { + log.Printf("[reconcile][退款自愈] out_trade_no=%s 本地曾卡于 %s,succeeded 退款 %d/%d → 重新 apply 为 %s", + outTradeNo, o.Status, succ, o.AmountMinor, next) + } + return nil +} + +// RefundStuckAlertTask 是「卡滞 processing/manual_pending 退款告警」义务:同 sweep +// 家族的姊妹任务,只读观测 —— 只打 WARN,绝不改状态(状态机翻转是 +// RefundApplyTask/业务方的事)。渠道退款查询 API 面(主动向渠道问退款进度)留待后续; +// 这里先用「本地卡滞时长」兜底可见性。 +func RefundStuckAlertTask(refunds *store.RefundStore, threshold time.Duration, now func() time.Time) func(ctx context.Context) error { + return func(ctx context.Context) error { + cutoff := now().Add(-threshold) + stuck, err := refunds.ListStuckRefunds( + []model.RefundStatus{model.RefundProcessing, model.RefundManualPending}, cutoff, 200) + if err != nil { + return err + } + for i := range stuck { + r := &stuck[i] + log.Printf("[reconcile][WARN][退款卡滞] refund_id=%s out_trade_no=%s status=%s 已卡滞 %s(阈值 %s)——"+ + "仅观测告警不改状态;渠道退款查询 API 面留待后续", + r.RefundID, r.OutTradeNo, r.Status, now().Sub(r.UpdatedAt).Round(time.Minute), threshold) + } + return nil + } +} diff --git a/internal/reconcile/refund_apply_test.go b/internal/reconcile/refund_apply_test.go new file mode 100644 index 0000000..8a3c752 --- /dev/null +++ b/internal/reconcile/refund_apply_test.go @@ -0,0 +1,189 @@ +package reconcile_test + +import ( + "bytes" + "context" + "log" + "strings" + "testing" + "time" + + "github.com/wangjia/pay/internal/model" + "github.com/wangjia/pay/internal/reconcile" + "github.com/wangjia/pay/internal/store" +) + +// TestRefundApplyTaskSelfHealsStuckPaidOrder 钉住 P4 T3 review 的崩溃窗口: +// MarkRefundStatus 把退款翻成 succeeded 后、调用方在调 ApplyRefundToOrder 前崩溃, +// 订单卡在 paid。RefundApplyTask 应重算 succeeded 之和并重新 apply,自愈成 +// refunded(全额)。 +func TestRefundApplyTaskSelfHealsStuckPaidOrder(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + refunds := store.NewRefundStore(db) + + if err := orders.CreateOrder(&model.OrderV2{OutTradeNo: "STUCK-1", AmountMinor: 10000, Currency: "CNY", Status: model.OrderPaidV2}); err != nil { + t.Fatal(err) + } + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-stuck-1", OutTradeNo: "STUCK-1", AttemptProviderRef: "STUCK-1", + AmountMinor: 10000, Currency: "CNY", Status: model.RefundSucceeded, InitiatedBy: "business"}); err != nil { + t.Fatal(err) + } + + task := reconcile.RefundApplyTask(orders, refunds, 200) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + + o, err := orders.GetOrder("STUCK-1") + if err != nil || o.Status != model.OrderRefundedV2 { + t.Fatalf("应自愈为 refunded, got %+v, err=%v", o, err) + } +} + +// TestRefundApplyTaskPartialStaysPartial 部分退款(succeeded 之和 < 订单金额)应自愈为 +// partially_refunded,而非误判 fully。 +func TestRefundApplyTaskPartialStaysPartial(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + refunds := store.NewRefundStore(db) + + if err := orders.CreateOrder(&model.OrderV2{OutTradeNo: "STUCK-2", AmountMinor: 10000, Currency: "CNY", Status: model.OrderPaidV2}); err != nil { + t.Fatal(err) + } + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-stuck-2", OutTradeNo: "STUCK-2", AttemptProviderRef: "STUCK-2", + AmountMinor: 4000, Currency: "CNY", Status: model.RefundSucceeded, InitiatedBy: "business"}); err != nil { + t.Fatal(err) + } + + task := reconcile.RefundApplyTask(orders, refunds, 200) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + o, err := orders.GetOrder("STUCK-2") + if err != nil || o.Status != model.OrderPartRefundedV2 { + t.Fatalf("应自愈为 partially_refunded, got %+v, err=%v", o, err) + } +} + +// TestRefundApplyTaskIdempotentNoOpOnRerun 幂等:自愈一次后重跑不应报错、不应再次 +// "翻转"(状态已一致,ApplyRefundToOrder 不应被重复触发出错误的副作用)。 +func TestRefundApplyTaskIdempotentNoOpOnRerun(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + refunds := store.NewRefundStore(db) + + if err := orders.CreateOrder(&model.OrderV2{OutTradeNo: "STUCK-3", AmountMinor: 5000, Currency: "CNY", Status: model.OrderPaidV2}); err != nil { + t.Fatal(err) + } + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-stuck-3", OutTradeNo: "STUCK-3", AttemptProviderRef: "STUCK-3", + AmountMinor: 5000, Currency: "CNY", Status: model.RefundSucceeded, InitiatedBy: "business"}); err != nil { + t.Fatal(err) + } + + task := reconcile.RefundApplyTask(orders, refunds, 200) + for i := 0; i < 3; i++ { + if err := task(context.Background()); err != nil { + t.Fatalf("run %d: %v", i, err) + } + } + o, err := orders.GetOrder("STUCK-3") + if err != nil || o.Status != model.OrderRefundedV2 { + t.Fatalf("重跑后仍应 refunded, got %+v, err=%v", o, err) + } +} + +// TestRefundApplyTaskConsistentOrderUntouched 已一致(无 succeeded 退款,或订单已是 +// 该退款对应的终态)的订单不应被误触发。 +func TestRefundApplyTaskConsistentOrderUntouched(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + refunds := store.NewRefundStore(db) + + if err := orders.CreateOrder(&model.OrderV2{OutTradeNo: "OK-1", AmountMinor: 1000, Currency: "CNY", Status: model.OrderPaidV2}); err != nil { + t.Fatal(err) + } + // 有一笔 processing(未 succeeded)退款,不该触发自愈。 + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-ok-1", OutTradeNo: "OK-1", AttemptProviderRef: "OK-1", + AmountMinor: 500, Currency: "CNY", Status: model.RefundProcessing, InitiatedBy: "business"}); err != nil { + t.Fatal(err) + } + + task := reconcile.RefundApplyTask(orders, refunds, 200) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + o, err := orders.GetOrder("OK-1") + if err != nil || o.Status != model.OrderPaidV2 { + t.Fatalf("无 succeeded 退款不应被翻转, got %+v, err=%v", o, err) + } +} + +// TestRefundStuckAlertTaskLogsWarnWithoutChangingState 卡滞 processing/manual_pending +// 超阈值只应打 WARN 日志,不改任何状态(观测型)。 +func TestRefundStuckAlertTaskLogsWarnWithoutChangingState(t *testing.T) { + db := model.OpenTestDB(t) + refunds := store.NewRefundStore(db) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-stall-1", OutTradeNo: "STALL-1", + AmountMinor: 100, Currency: "CNY", Status: model.RefundProcessing}); err != nil { + t.Fatal(err) + } + if err := db.Model(&model.Refund{}).Where("refund_id = ?", "rf-stall-1"). + Update("updated_at", now.Add(-45*time.Minute)).Error; err != nil { + t.Fatal(err) + } + + var buf bytes.Buffer + orig := log.Writer() + log.SetOutput(&buf) + defer log.SetOutput(orig) + + task := reconcile.RefundStuckAlertTask(refunds, 30*time.Minute, func() time.Time { return now }) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + + if !strings.Contains(buf.String(), "rf-stall-1") || !strings.Contains(buf.String(), "WARN") { + t.Fatalf("应打 WARN 日志含 refund_id, got: %s", buf.String()) + } + + r, err := refunds.GetRefund("rf-stall-1") + if err != nil || r.Status != model.RefundProcessing { + t.Fatalf("告警不应改状态, got %+v, err=%v", r, err) + } +} + +// TestRefundStuckAlertTaskSkipsFreshAndTerminal 未超阈值 / 已终态的退款不应被告警。 +func TestRefundStuckAlertTaskSkipsFreshAndTerminal(t *testing.T) { + db := model.OpenTestDB(t) + refunds := store.NewRefundStore(db) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-fresh", OutTradeNo: "FRESH-1", + AmountMinor: 100, Currency: "CNY", Status: model.RefundProcessing}); err != nil { + t.Fatal(err) + } + if err := refunds.CreateRefund(&model.Refund{RefundID: "rf-done", OutTradeNo: "DONE-1", + AmountMinor: 100, Currency: "CNY", Status: model.RefundSucceeded}); err != nil { + t.Fatal(err) + } + if err := db.Model(&model.Refund{}).Where("refund_id = ?", "rf-done"). + Update("updated_at", now.Add(-2*time.Hour)).Error; err != nil { + t.Fatal(err) + } + + var buf bytes.Buffer + orig := log.Writer() + log.SetOutput(&buf) + defer log.SetOutput(orig) + + task := reconcile.RefundStuckAlertTask(refunds, 30*time.Minute, func() time.Time { return now }) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + if strings.Contains(buf.String(), "rf-fresh") || strings.Contains(buf.String(), "rf-done") { + t.Fatalf("未超阈值/已终态不应告警, got: %s", buf.String()) + } +} diff --git a/internal/reconcile/sync.go b/internal/reconcile/sync.go new file mode 100644 index 0000000..749b2a4 --- /dev/null +++ b/internal/reconcile/sync.go @@ -0,0 +1,59 @@ +package reconcile + +import ( + "context" + "log" + "time" + + "github.com/wangjia/pay/internal/gateway" + "github.com/wangjia/pay/internal/provider" + "github.com/wangjia/pay/internal/store" +) + +// SyncPendingTask 调度 P2 gateway.SyncPendingAttempts:逐 pending attempt 查单收敛(防掉单)。 +// 其内部已对 not_found/amount_mismatch/failed 打日志(settle-sync,764ed55),此处不重复。 +func SyncPendingTask(gw *gateway.Gateway, limit int) func(ctx context.Context) error { + return func(ctx context.Context) error { + n, err := gw.SyncPendingAttempts(ctx, limit) + if err != nil { + return err + } + if n > 0 { + log.Printf("[reconcile] 查单对账收敛 %d 笔待支付 → paid", n) + } + return nil + } +} + +// PaidSpotCheckTask 已付订单抽查:对近 window 内已付 attempt 反查渠道,金额/币种漂移即告警 +// (如渠道侧已退款/拒付而本地仍 paid)。只发现不改状态——状态机翻转属 P4。 +// crypto 之类 query-only 渠道:paid 后再查若命中同额即一致;查不到(链上历史滚出窗口)不报错跳过。 +func PaidSpotCheckTask(orders *store.OrderStore, providers *provider.Registry, window time.Duration, now func() time.Time) func(ctx context.Context) error { + return func(ctx context.Context) error { + atts, err := orders.ListRecentlyPaidAttempts(now().Add(-window), 100) + if err != nil { + return err + } + for i := range atts { + a := &atts[i] + prov, err := providers.Get(a.Channel) + if err != nil { + continue + } + created := a.CreatedAt + ev, err := prov.Query(ctx, provider.QueryRequest{ + ProviderRef: a.ProviderRef, OutTradeNo: a.OutTradeNo, AccountID: a.AccountID, + AmountMinor: a.AmountMinor, Currency: a.Currency, CreatedAt: created, ExpiresAt: a.ExpiresAt, + }) + if err != nil || ev == nil { + continue // 查不到/瞬时错:抽查尽力而为,不阻断 + } + // 本地 paid,渠道却报非成功,或金额/币种对不上 → 对账差异,必须可见。 + if ev.Status != provider.PaidSucceeded || ev.PaidCurrency != a.Currency || ev.PaidAmountMinor < a.AmountMinor { + log.Printf("[reconcile][对账差异] attempt=%s channel=%s 本地 paid 但渠道 status=%s amount=%d/%s(本地 %d/%s)", + a.ProviderRef, a.Channel, ev.Status, ev.PaidAmountMinor, ev.PaidCurrency, a.AmountMinor, a.Currency) + } + } + return nil + } +} diff --git a/internal/reconcile/sync_test.go b/internal/reconcile/sync_test.go new file mode 100644 index 0000000..475aca6 --- /dev/null +++ b/internal/reconcile/sync_test.go @@ -0,0 +1,72 @@ +package reconcile_test + +import ( + "context" + "testing" + "time" + + "github.com/wangjia/pay/config" + "github.com/wangjia/pay/internal/accounts" + "github.com/wangjia/pay/internal/gateway" + "github.com/wangjia/pay/internal/model" + "github.com/wangjia/pay/internal/provider" + "github.com/wangjia/pay/internal/provider/fake" + "github.com/wangjia/pay/internal/reconcile" + "github.com/wangjia/pay/internal/store" +) + +type stubResolver struct{} + +func (stubResolver) Resolve(sku, currency string) (int64, string, string, error) { + return 29990000, "Pro", "pro_year", nil +} + +type nopEnq struct{} + +// Enqueue 签名依当前仓库 gateway.WebhookEnqueuer(4 个 string + map;brief 草稿只写 3 个 — +// P4 T2 之后加了 refundID 参数,此处适配现状,见 task-5-report.md 记录的漂移)。 +func (nopEnq) Enqueue(_, _, _, _ string, _ map[string]any) error { return nil } + +func TestSyncPendingTaskSettlesViaQuery(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + preg := provider.NewRegistry() + fp := fake.New() + preg.Register(fp) + areg := accounts.New([]config.AccountConfig{{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true}}) + gw := gateway.New(orders, preg, accounts.NewRouter(areg, nil, nil), stubResolver{}, nopEnq{}, "global") + + res, _ := gw.CreateOrder(context.Background(), gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"}) + atts, _ := orders.ListAttemptsByStatus(model.AttemptPending, 10) + fp.SetQueryResult(atts[0].ProviderRef, provider.PaidEvent{ + ProviderRef: atts[0].ProviderRef, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"}) + + task := reconcile.SyncPendingTask(gw, 50) + if err := task(context.Background()); err != nil { + t.Fatalf("task: %v", err) + } + o, _ := orders.GetOrder(res.OrderNo) + if o.Status != model.OrderPaidV2 { + t.Fatalf("查单对账后应 paid, got %v", o.Status) + } + _ = time.Second +} + +func TestPaidSpotCheckTaskRunsCleanOnConsistent(t *testing.T) { + db := model.OpenTestDB(t) + orders := store.NewOrderStore(db) + preg := provider.NewRegistry() + fp := fake.New() + preg.Register(fp) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + paid := now.Add(-10 * time.Minute) + _ = orders.CreateAttempt(&model.Attempt{OutTradeNo: "O1", Channel: "fake", ProviderRef: "R-O1", + AmountMinor: 100, Currency: "USDT", Status: model.AttemptPaid, PaidAt: &paid}) + // 渠道侧查单仍报 succeeded 同额 → 一致,无告警。 + fp.SetQueryResult("R-O1", provider.PaidEvent{ProviderRef: "R-O1", Status: provider.PaidSucceeded, PaidAmountMinor: 100, PaidCurrency: "USDT"}) + + task := reconcile.PaidSpotCheckTask(orders, preg, time.Hour, func() time.Time { return now }) + if err := task(context.Background()); err != nil { + t.Fatalf("spotcheck: %v", err) // 只求不报错、不 panic;漂移检测走日志 + } +} diff --git a/internal/store/order_query.go b/internal/store/order_query.go index 53884d0..1f07bae 100644 --- a/internal/store/order_query.go +++ b/internal/store/order_query.go @@ -147,3 +147,30 @@ func (s *OrderStore) SumPaidAttemptMinorByAccountSince(since time.Time) (map[str } return out, nil } + +// ListRecentlyPaidAttempts 列近期(paid_at>=since)已付 attempt,供对账抽查反查渠道核对。 +func (s *OrderStore) ListRecentlyPaidAttempts(since time.Time, limit int) ([]model.Attempt, error) { + if limit <= 0 || limit > 500 { + limit = 100 + } + var out []model.Attempt + if err := s.db.Where("status = ? AND paid_at >= ?", model.AttemptPaid, since). + Order("id DESC").Limit(limit).Find(&out).Error; err != nil { + return nil, fmt.Errorf("store.ListRecentlyPaidAttempts: %w", err) + } + return out, nil +} + +// ListOrdersByStatus 按状态集合列订单,供退款修复扫描(Task 5 P4 义务)定位「当前处于 +// 退款相关态」的候选订单,与 RefundStore.ListDistinctOutTradeNosByStatus(succeeded)取并集。 +func (s *OrderStore) ListOrdersByStatus(statuses []model.OrderStatusV2, limit int) ([]model.OrderV2, error) { + if limit <= 0 || limit > 500 { + limit = 200 + } + var out []model.OrderV2 + if err := s.db.Where("status IN ?", statuses). + Order("id ASC").Limit(limit).Find(&out).Error; err != nil { + return nil, fmt.Errorf("store.ListOrdersByStatus: %w", err) + } + return out, nil +} diff --git a/internal/store/order_query_test.go b/internal/store/order_query_test.go index 820714d..968aa50 100644 --- a/internal/store/order_query_test.go +++ b/internal/store/order_query_test.go @@ -123,3 +123,56 @@ func TestSumPaidAttemptMinorByAccountSince(t *testing.T) { t.Fatalf("聚合 = %+v, want acct-1=15000 acct-2=7000", got) } } + +func TestListRecentlyPaidAttempts(t *testing.T) { + db := model.OpenTestDB(t) + s := store.NewOrderStore(db) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + mk := func(no string, st model.AttemptStatus, paidAgo time.Duration) { + paid := now.Add(-paidAgo) + _ = s.CreateAttempt(&model.Attempt{OutTradeNo: no, Channel: "fake", ProviderRef: "R-" + no, + AmountMinor: 100, Currency: "USDT", Status: st, PaidAt: &paid}) + } + mk("RECENT", model.AttemptPaid, 10*time.Minute) // 近期已付 → 命中 + mk("OLD", model.AttemptPaid, 5*time.Hour) // 太旧 → 不命中 + mk("PEND", model.AttemptPending, 1*time.Minute) // 未付 → 不命中 + + got, err := s.ListRecentlyPaidAttempts(now.Add(-time.Hour), 50) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(got) != 1 || got[0].OutTradeNo != "RECENT" { + t.Fatalf("只应含 RECENT, got %+v", got) + } +} + +// TestListOrdersByStatus 供 P6+P4 义务的「退款修复扫描」定位候选订单:按状态集合 +// 列订单(如 refunding/partially_refunded),与 refund 表的 succeeded 记录取并集 +// 作为重算候选。 +func TestListOrdersByStatus(t *testing.T) { + db := model.OpenTestDB(t) + s := store.NewOrderStore(db) + mk := func(no string, st model.OrderStatusV2) { + _ = s.CreateOrder(&model.OrderV2{OutTradeNo: no, AmountMinor: 100, Currency: "CNY", Status: st}) + } + mk("O-PAID", model.OrderPaidV2) + mk("O-REFUNDING", model.OrderRefundingV2) + mk("O-PART", model.OrderPartRefundedV2) + mk("O-DONE", model.OrderRefundedV2) + mk("O-PENDING", model.OrderPendingV2) + + got, err := s.ListOrdersByStatus([]model.OrderStatusV2{model.OrderRefundingV2, model.OrderPartRefundedV2}, 50) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(got) != 2 { + t.Fatalf("应命中 2 张(REFUNDING/PART), got %d: %+v", len(got), got) + } + seen := map[string]bool{} + for _, o := range got { + seen[o.OutTradeNo] = true + } + if !seen["O-REFUNDING"] || !seen["O-PART"] { + t.Fatalf("命中集合不对: %+v", got) + } +} diff --git a/internal/store/refund.go b/internal/store/refund.go index 4f93df4..49370be 100644 --- a/internal/store/refund.go +++ b/internal/store/refund.go @@ -128,6 +128,38 @@ func (s *RefundStore) MarkRefundStatus(refundID string, from, to model.RefundSta return res.RowsAffected > 0, nil } +// ListDistinctOutTradeNosByStatus 列有某状态退款的 distinct out_trade_no,供退款修复 +// 扫描(RefundApplyTask)定位「有 succeeded 退款」的候选订单——self-heal「退款成功但 +// 订单卡 paid」的崩溃窗口(P4 T3 review 义务,见 task-5-brief 外的两条追加义务)。 +func (s *RefundStore) ListDistinctOutTradeNosByStatus(status model.RefundStatus, limit int) ([]string, error) { + if limit <= 0 || limit > 500 { + limit = 200 + } + var out []string + if err := s.db.Model(&model.Refund{}).Where("status = ?", status). + Group("out_trade_no").Order("out_trade_no ASC").Limit(limit). + Pluck("out_trade_no", &out).Error; err != nil { + return nil, fmt.Errorf("store.ListDistinctOutTradeNosByStatus: %w", err) + } + return out, nil +} + +// ListStuckRefunds 列 status 落在给定集合、且 updated_at 早于 before 的退款行,供 +// 「卡滞 processing/manual_pending 退款告警」只读观测扫描用(不改状态)。updated_at +// 用作「进入当前状态」的近似时刻——本表除 MarkRefundStatus/CreateRefundGuarded 外 +// 不写,近似成立。 +func (s *RefundStore) ListStuckRefunds(statuses []model.RefundStatus, before time.Time, limit int) ([]model.Refund, error) { + if limit <= 0 || limit > 500 { + limit = 200 + } + var out []model.Refund + if err := s.db.Where("status IN ? AND updated_at < ?", statuses, before). + Order("updated_at ASC").Limit(limit).Find(&out).Error; err != nil { + return nil, fmt.Errorf("store.ListStuckRefunds: %w", err) + } + return out, nil +} + // ListManualPending lists refunds awaiting manual (crypto) settlement. func (s *RefundStore) ListManualPending(limit int) ([]model.Refund, error) { if limit <= 0 || limit > 200 { diff --git a/internal/store/refund_test.go b/internal/store/refund_test.go index 906f925..7370a5e 100644 --- a/internal/store/refund_test.go +++ b/internal/store/refund_test.go @@ -100,6 +100,68 @@ func TestApplyRefundToOrderFully(t *testing.T) { } } +// TestListDistinctOutTradeNosByStatus 供退款修复扫描定位「有 succeeded 退款」的候选 +// 订单(自愈依据):同订单多笔 succeeded 退款只应出现一次(distinct)。 +func TestListDistinctOutTradeNosByStatus(t *testing.T) { + db := model.OpenTestDB(t) + rs := NewRefundStore(db) + _ = rs.CreateRefund(&model.Refund{RefundID: "s1", OutTradeNo: "PAY-S1", AmountMinor: 100, Currency: "CNY", Status: model.RefundSucceeded}) + _ = rs.CreateRefund(&model.Refund{RefundID: "s2", OutTradeNo: "PAY-S1", AmountMinor: 200, Currency: "CNY", Status: model.RefundSucceeded}) // 同单第二笔 + _ = rs.CreateRefund(&model.Refund{RefundID: "s3", OutTradeNo: "PAY-S2", AmountMinor: 100, Currency: "CNY", Status: model.RefundSucceeded}) + _ = rs.CreateRefund(&model.Refund{RefundID: "p1", OutTradeNo: "PAY-S3", AmountMinor: 100, Currency: "CNY", Status: model.RefundProcessing}) // 非 succeeded,不应命中 + + got, err := rs.ListDistinctOutTradeNosByStatus(model.RefundSucceeded, 50) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(got) != 2 { + t.Fatalf("应 distinct 出 2 个 out_trade_no, got %d: %+v", len(got), got) + } + seen := map[string]bool{} + for _, no := range got { + seen[no] = true + } + if !seen["PAY-S1"] || !seen["PAY-S2"] { + t.Fatalf("命中集合不对: %+v", got) + } +} + +// TestListStuckRefunds 供「卡滞 processing/manual_pending 退款告警」:只挑 updated_at +// 早于阈值的 processing/manual_pending 行,requested/succeeded/failed 不命中。 +func TestListStuckRefunds(t *testing.T) { + db := model.OpenTestDB(t) + rs := NewRefundStore(db) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + mk := func(id, no string, st model.RefundStatus, updatedAgo time.Duration) { + if err := rs.CreateRefund(&model.Refund{RefundID: id, OutTradeNo: no, AmountMinor: 100, Currency: "CNY", Status: st}); err != nil { + t.Fatal(err) + } + if err := db.Model(&model.Refund{}).Where("refund_id = ?", id). + Update("updated_at", now.Add(-updatedAgo)).Error; err != nil { + t.Fatal(err) + } + } + mk("stuck-proc", "PAY-T1", model.RefundProcessing, 45*time.Minute) // 超阈值 → 命中 + mk("fresh-proc", "PAY-T2", model.RefundProcessing, 5*time.Minute) // 未超 → 不命中 + mk("stuck-manual", "PAY-T3", model.RefundManualPending, 2*time.Hour) // 超阈值 → 命中 + mk("done", "PAY-T4", model.RefundSucceeded, 2*time.Hour) // 已终态 → 不命中 + + got, err := rs.ListStuckRefunds([]model.RefundStatus{model.RefundProcessing, model.RefundManualPending}, now.Add(-30*time.Minute), 50) + if err != nil { + t.Fatalf("list: %v", err) + } + if len(got) != 2 { + t.Fatalf("应命中 2 笔卡滞, got %d: %+v", len(got), got) + } + seen := map[string]bool{} + for _, r := range got { + seen[r.RefundID] = true + } + if !seen["stuck-proc"] || !seen["stuck-manual"] { + t.Fatalf("命中集合不对: %+v", got) + } +} + func TestListManualPending(t *testing.T) { db := model.OpenTestDB(t) rs := NewRefundStore(db) diff --git a/main.go b/main.go index 17813e7..2cd8cdc 100644 --- a/main.go +++ b/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "errors" "log" "strings" @@ -17,6 +18,7 @@ import ( "github.com/wangjia/pay/internal/gateway" "github.com/wangjia/pay/internal/model" "github.com/wangjia/pay/internal/providerbuild" + "github.com/wangjia/pay/internal/reconcile" "github.com/wangjia/pay/internal/router" "github.com/wangjia/pay/internal/store" "github.com/wangjia/pay/internal/webhook" @@ -53,11 +55,39 @@ func main() { acctReg := accounts.New(config.C.Accounts) pReg := providerbuild.BuildRegistry(acctReg) // 配置驱动:有 enabled 账户才注册对应渠道(P3) // P5 多账户路由:按 config.routing. 选策略(缺省 round_robin)。 - // limit_aware 用量数据源 P6 对账就绪前用空源(NopUsage,退化为 round_robin)。 - acctPicker := accounts.NewRouter(acctReg, config.C.Routing, accounts.NopUsage{}) + // limit_aware 用量数据源:P6 对账 Runner 周期 Refresh 的真实用量源(替 NopUsage)。 + usage := reconcile.NewUsageSource(orderStore, time.Now) + acctPicker := accounts.NewRouter(acctReg, config.C.Routing, usage) gw := gateway.New(orderStore, pReg, acctPicker, productResolver, notifier, "cn") router.SetupV2(r, gw) + // P6 后台守护 / 对账:订单过期清理 + 用量刷新 + 查单对账收敛 + 已付抽查 + + // 退款修复扫描/卡滞告警(P4 T3 review 追加义务,归到本任务一起装配)。 + // crypto 预留冷启动 Warm + 孤儿扫描留给 Task 6 追加(AddCryptoJobs)。 + if config.C.Reconcile.Enabled { + rc := config.C.Reconcile + refundStore := store.NewRefundStore(db) + runner := reconcile.NewRunner() + runner.Add("order-expire", time.Duration(rc.ExpireEverySec)*time.Second, + reconcile.OrderExpirerTask(orderStore, time.Duration(rc.OrderTTLMin)*time.Minute, time.Now)) + runner.Add("usage-refresh", time.Duration(rc.UsageEverySec)*time.Second, + reconcile.RefreshUsageTask(usage)) + runner.Add("sync-pending", time.Duration(rc.SyncEverySec)*time.Second, + reconcile.SyncPendingTask(gw, 100)) + runner.Add("paid-spotcheck", time.Duration(rc.SpotCheckEverySec)*time.Second, + reconcile.PaidSpotCheckTask(orderStore, pReg, time.Duration(rc.SpotCheckWindowMin)*time.Minute, time.Now)) + runner.Add("refund-apply-sweep", time.Duration(rc.RefundApplyEverySec)*time.Second, + reconcile.RefundApplyTask(orderStore, refundStore, 200)) + runner.Add("refund-stuck-alert", time.Duration(rc.RefundApplyEverySec)*time.Second, + reconcile.RefundStuckAlertTask(refundStore, time.Duration(rc.RefundStuckWarnMin)*time.Minute, time.Now)) + // crypto 孤儿扫描(Task 6)在此追加:reconcile.AddCryptoJobs(runner, pReg, acctReg, orderStore, rc) + + ctx := context.Background() + runner.RunOnce(ctx) // 启动预热:先跑一遍(usage 快照/过期清理/退款自愈立即生效) + runner.Start(ctx) + log.Printf("[reconcile] 后台守护已启动(过期清理/用量刷新/查单对账/已付抽查/退款修复扫描/退款卡滞告警)") + } + if config.C.QuerySync.Enabled { orderSvc.StartQuerySync( time.Duration(config.C.QuerySync.IntervalSec)*time.Second,