From e04cd73983f7a1f3c9ada081f70585698162f889 Mon Sep 17 00:00:00 2001 From: wangjia <809946525@qq.com> Date: Fri, 10 Jul 2026 17:29:08 +0800 Subject: [PATCH] =?UTF-8?q?feat(v2):=20crypto=20=E9=A2=84=E7=95=99?= =?UTF-8?q?=E8=A1=A8=E5=86=B7=E5=90=AF=E5=8A=A8=E5=85=9C=E5=BA=95=E2=80=94?= =?UTF-8?q?=E2=80=94=E6=B3=A8=E5=85=A5=20loader=20=E4=BB=8E=20pending=20at?= =?UTF-8?q?tempts=20=E9=87=8D=E5=BB=BA=20reservation(=E5=85=9C=E9=87=8D?= =?UTF-8?q?=E5=90=AF=E4=B8=A2=E5=86=85=E5=AD=98)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/provider/crypto/crypto.go | 52 ++++++++++++++++++++++++- internal/provider/crypto/crypto_test.go | 32 +++++++++++++++ internal/reconcile/crypto_warm.go | 33 ++++++++++++++++ internal/reconcile/crypto_warm_test.go | 29 ++++++++++++++ 4 files changed, 144 insertions(+), 2 deletions(-) create mode 100644 internal/reconcile/crypto_warm.go create mode 100644 internal/reconcile/crypto_warm_test.go diff --git a/internal/provider/crypto/crypto.go b/internal/provider/crypto/crypto.go index a6551cb..33b7520 100644 --- a/internal/provider/crypto/crypto.go +++ b/internal/provider/crypto/crypto.go @@ -49,15 +49,29 @@ type Provider struct { baseURL string http *http.Client now func() time.Time + loader ReservationLoader // 冷启动预留重建源(装配期注入,nil=不重建) mu sync.Mutex reserved map[string]time.Time // "
/" → 预留到期(链上匹配维度,对齐 Query 的 to==addr)(canonical AmountRecentlyUsed 的进程内等价) } +// PendingReservation 冷启动重建一笔预留所需的最小信息(中性结构,crypto 不依赖 store)。 +type PendingReservation struct { + AccountID string // 收款账户(用于解析地址,链上匹配维度) + AmountMinor int64 // attempt 冻结的 base 金额(不含尾数) + ProviderRef string // "CRYPTO--",用于恢复尾数 + ReservedAt time.Time // 建单时间(= attempt.CreatedAt),冷却窗自此算 +} + +// ReservationLoader 返回当前仍活跃(pending)的 crypto 预留。装配期由 main 用 OrderStore 实现。 +type ReservationLoader func(ctx context.Context) ([]PendingReservation, error) + type Option func(*Provider) -func WithBaseURL(u string) Option { return func(p *Provider) { p.baseURL = u } } -func WithHTTPClient(c *http.Client) Option { return func(p *Provider) { p.http = c } } +func WithBaseURL(u string) Option { return func(p *Provider) { p.baseURL = u } } +func WithHTTPClient(c *http.Client) Option { return func(p *Provider) { p.http = c } } +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 } } func New(accts *accounts.Registry, opts ...Option) *Provider { p := &Provider{ @@ -151,6 +165,40 @@ func tailFromRef(ref string) (int64, error) { return strconv.ParseInt(ref[i+1:], 10, 64) } +// Warm 冷启动兜底:把仍在冷却窗内的活跃预留灌回内存表,兜住重启丢 map 导致的金额复用误配。 +// 幂等:只加不覆盖更早到期时间;冷却已过的跳过。装配期在起服务前调一次即可。 +func (p *Provider) Warm(ctx context.Context) error { + if p.loader == nil { + return nil + } + items, err := p.loader(ctx) + if err != nil { + return err + } + now := p.now() + p.mu.Lock() + defer p.mu.Unlock() + for _, it := range items { + addr, err := p.address(it.AccountID) // 地址是链上匹配维度真相源 + if err != nil { + continue + } + tail, err := tailFromRef(it.ProviderRef) + if err != nil { + continue + } + until := it.ReservedAt.Add(amountCooldown) + if !until.After(now) { + continue // 冷却已过,金额可安全复用,无需恢复 + } + key := addr + "/" + strconv.FormatInt(it.AmountMinor+tail, 10) + if cur, ok := p.reserved[key]; !ok || until.After(cur) { + p.reserved[key] = until + } + } + return nil +} + func (p *Provider) Create(_ context.Context, req provider.CreateRequest) (*provider.Session, error) { if req.Currency != "USDT" { return nil, fmt.Errorf("crypto: 仅支持 USDT, got %s", req.Currency) diff --git a/internal/provider/crypto/crypto_test.go b/internal/provider/crypto/crypto_test.go index 0be18aa..d6fdcea 100644 --- a/internal/provider/crypto/crypto_test.go +++ b/internal/provider/crypto/crypto_test.go @@ -206,3 +206,35 @@ func TestReservationKeyByAddress(t *testing.T) { t.Fatalf("共享地址的两个账户不应分配相同金额: %d", amt1) } } + +// 冷启动兜底:装配期注入的 ReservationLoader 在 Warm 时把仍在冷却窗内的 pending +// 预留灌回内存表;超冷却窗的旧预留(迟到旧款已不可能匹配)不必恢复。 +func TestWarmRebuildsReservationsFromLoader(t *testing.T) { + const addr = "TWarmTestAddr000000000000000000000" + t.Setenv("CRY_ADDRESS", addr) + reg := accounts.New([]config.AccountConfig{ + {AccountID: "cry-1", Channel: "crypto", Enabled: true, CredentialEnvPrefix: "cry"}, + }) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + + // 两条 pending:一条在冷却窗内(应恢复),一条建单于 40min 前(> 30min 冷却窗,应跳过)。 + loader := func(context.Context) ([]crypto.PendingReservation, error) { + return []crypto.PendingReservation{ + {AccountID: "cry-1", AmountMinor: 29990000, ProviderRef: "CRYPTO-PAY-A-263", ReservedAt: now.Add(-5 * time.Minute)}, + {AccountID: "cry-1", AmountMinor: 29990000, ProviderRef: "CRYPTO-PAY-B-777", ReservedAt: now.Add(-40 * time.Minute)}, + }, nil + } + p := crypto.New(reg, crypto.WithReservationLoader(loader), crypto.WithNow(func() time.Time { return now })) + if err := p.Warm(context.Background()); err != nil { + t.Fatalf("warm: %v", err) + } + res := p.GetReserved() + inWindow := addr + "/" + strconv.FormatInt(29990000+263, 10) + expired := addr + "/" + strconv.FormatInt(29990000+777, 10) + if _, ok := res[inWindow]; !ok { + t.Fatalf("冷却窗内的预留应恢复, got %v", res) + } + if _, ok := res[expired]; ok { + t.Fatalf("超冷却窗的预留不应恢复, got %v", res) + } +} diff --git a/internal/reconcile/crypto_warm.go b/internal/reconcile/crypto_warm.go new file mode 100644 index 0000000..f62b69a --- /dev/null +++ b/internal/reconcile/crypto_warm.go @@ -0,0 +1,33 @@ +package reconcile + +import ( + "context" + + "github.com/wangjia/pay/internal/model" + "github.com/wangjia/pay/internal/provider/crypto" + "github.com/wangjia/pay/internal/store" +) + +// CryptoReservationLoader 装配 crypto 冷启动预留源:读全部 pending attempt, +// 过滤 channel=crypto,映射为 crypto.PendingReservation。ReservedAt 取 attempt.CreatedAt +// (建单时刻,冷却窗自此算)。crypto 不 import store,故此桥在装配层。 +func CryptoReservationLoader(orders *store.OrderStore) crypto.ReservationLoader { + return func(ctx context.Context) ([]crypto.PendingReservation, error) { + atts, err := orders.ListAttemptsByStatus(model.AttemptPending, 200) + if err != nil { + return nil, err + } + out := make([]crypto.PendingReservation, 0, len(atts)) + for i := range atts { + a := &atts[i] + if a.Channel != "crypto" { + continue + } + out = append(out, crypto.PendingReservation{ + AccountID: a.AccountID, AmountMinor: a.AmountMinor, + ProviderRef: a.ProviderRef, ReservedAt: a.CreatedAt, + }) + } + return out, nil + } +} diff --git a/internal/reconcile/crypto_warm_test.go b/internal/reconcile/crypto_warm_test.go new file mode 100644 index 0000000..95baffc --- /dev/null +++ b/internal/reconcile/crypto_warm_test.go @@ -0,0 +1,29 @@ +package reconcile_test + +import ( + "context" + "testing" + + "github.com/wangjia/pay/internal/model" + "github.com/wangjia/pay/internal/reconcile" + "github.com/wangjia/pay/internal/store" +) + +func TestCryptoReservationLoaderFiltersPendingCrypto(t *testing.T) { + db := model.OpenTestDB(t) + s := store.NewOrderStore(db) + // 两条 attempt:crypto pending(要)、alipay pending(不要)。 + _ = s.CreateAttempt(&model.Attempt{OutTradeNo: "O1", Channel: "crypto", ProviderRef: "CRYPTO-O1-12", + AmountMinor: 100, Currency: "USDT", Status: model.AttemptPending}) + _ = s.CreateAttempt(&model.Attempt{OutTradeNo: "O2", Channel: "alipay", ProviderRef: "AL-O2", + AmountMinor: 200, Currency: "CNY", Status: model.AttemptPending}) + + loader := reconcile.CryptoReservationLoader(s) + items, err := loader(context.Background()) + if err != nil { + t.Fatalf("loader: %v", err) + } + if len(items) != 1 || items[0].ProviderRef != "CRYPTO-O1-12" || items[0].AmountMinor != 100 { + t.Fatalf("只应含 crypto pending, got %+v", items) + } +}