merge: P6 对账守护并入(七job装配/退款自愈/孤儿发现/Notifier硬化/真实用量源)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013nMthbVEmQquxBRKb9Fj8u

# Conflicts:
#	main.go
This commit is contained in:
wangjia
2026-07-10 18:16:27 +08:00
32 changed files with 2092 additions and 22 deletions
+117 -2
View File
@@ -49,15 +49,33 @@ type Provider struct {
baseURL string
http *http.Client
now func() time.Time
loader ReservationLoader // 冷启动预留重建源(装配期注入,nil=不重建)
mu sync.Mutex
reserved map[string]time.Time // "<address>/<amount>" → 预留到期(链上匹配维度,对齐 Query 的 to==addr)(canonical AmountRecentlyUsed 的进程内等价)
}
// PendingReservation 冷启动重建一笔预留所需的最小信息(中性结构,crypto 不依赖 store)。
type PendingReservation struct {
AccountID string // 收款账户(用于解析地址,链上匹配维度)
AmountMinor int64 // attempt 冻结的 base 金额(不含尾数)
ProviderRef string // "CRYPTO-<OutTradeNo>-<tail>",用于恢复尾数
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 } }
// 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{
@@ -151,6 +169,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)
@@ -188,6 +240,69 @@ func (p *Provider) VerifyCallback(_ context.Context, _ provider.CallbackInput) (
return nil, provider.ErrNotSupported
}
// ScanOrphans 扫地址近 Since 的确认到账,金额不在"任一 Known 的期望金额集"内 → 孤儿。
// 期望金额 = known.AmountMinor + tailFromRef(known.ProviderRef);块时须晚于 Since。
func (p *Provider) ScanOrphans(ctx context.Context, req provider.OrphanScanRequest) ([]provider.OrphanTransfer, error) {
addr, err := p.address(req.AccountID)
if err != nil {
return nil, err
}
expected := make(map[int64]struct{}, len(req.Known))
for _, k := range req.Known {
tail, terr := tailFromRef(k.ProviderRef)
if terr != nil {
continue // 无尾数的 ref 跳过(不误判为孤儿依据)
}
expected[k.AmountMinor+tail] = struct{}{}
}
endpoint := fmt.Sprintf("%s/v1/accounts/%s/transactions/trc20?only_confirmed=true&contract_address=%s&limit=50",
p.baseURL, url.PathEscape(addr), url.QueryEscape(USDTContract))
httpReq, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil)
if err != nil {
return nil, err
}
if k := p.apiKey(req.AccountID); k != "" {
httpReq.Header.Set("TRON-PRO-API-KEY", k)
}
resp, err := p.http.Do(httpReq)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("crypto: TronGrid HTTP %d: %s", resp.StatusCode, body)
}
var tr trc20Resp
if err := json.Unmarshal(body, &tr); err != nil {
return nil, err
}
sinceUnix := req.Since.Unix()
var out []provider.OrphanTransfer
for _, d := range tr.Data {
if d.To != addr || d.Type != "Transfer" {
continue
}
blockTs := d.BlockMs / 1000
if blockTs < sinceUnix {
continue // 窗外旧款不扫(避免把历史正常单反复报孤儿)
}
val, perr := strconv.ParseInt(d.Value, 10, 64)
if perr != nil {
continue
}
if _, ok := expected[val]; ok {
continue // 金额有主(匹配某 attempt 期望额)→ 非孤儿
}
out = append(out, provider.OrphanTransfer{
TxID: d.TxID, AmountMinor: val, Currency: "USDT", At: time.Unix(blockTs, 0),
})
}
return out, nil
}
// trc20Resp 对应 TronGrid /v1/accounts/{addr}/transactions/trc20 响应
// (canonical tron/client.go 同构;contract_address 查询参数已在服务端过滤合约)。
type trc20Resp struct {
+63
View File
@@ -206,3 +206,66 @@ 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)
}
}
// 孤儿到账发现:假 TronGrid 返回两笔确认到账,一笔匹配 known 期望金额(base+tail),
// 一笔无主 → 仅后者报孤儿。
func TestScanOrphansFlagsUnmatchedTransfer(t *testing.T) {
const addr = "TOrphanScanAddr00000000000000000000"
now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
// 假 TronGrid:to=addr 两笔确认到账。29990263 匹配 known(base 29990000 + tail 263);88880000 无主。
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"data":[
{"transaction_id":"TX-MATCH","to":"` + addr + `","type":"Transfer","value":"29990263","block_timestamp":` + strconv.FormatInt(now.Add(-5*time.Minute).UnixMilli(), 10) + `},
{"transaction_id":"TX-ORPHAN","to":"` + addr + `","type":"Transfer","value":"88880000","block_timestamp":` + strconv.FormatInt(now.Add(-3*time.Minute).UnixMilli(), 10) + `}
]}`))
}))
defer ts.Close()
t.Setenv("CRY_ADDRESS", addr)
reg := accounts.New([]config.AccountConfig{{AccountID: "cry-1", Channel: "crypto", Enabled: true, CredentialEnvPrefix: "cry"}})
p := crypto.New(reg, crypto.WithBaseURL(ts.URL), crypto.WithHTTPClient(ts.Client()), crypto.WithNow(func() time.Time { return now }))
orphans, err := p.ScanOrphans(context.Background(), provider.OrphanScanRequest{
AccountID: "cry-1", Since: now.Add(-time.Hour),
Known: []provider.KnownAttempt{{AmountMinor: 29990000, ProviderRef: "CRYPTO-PAY-A-263"}}, // 期望 29990263
})
if err != nil {
t.Fatalf("scan: %v", err)
}
if len(orphans) != 1 || orphans[0].TxID != "TX-ORPHAN" || orphans[0].AmountMinor != 88880000 {
t.Fatalf("只应报 1 笔孤儿 TX-ORPHAN, got %+v", orphans)
}
}
+30
View File
@@ -133,6 +133,36 @@ type RecurringProvider interface {
CancelAgreement(ctx context.Context, agreementRef string) error
}
// ---- 对账:孤儿到账扫描(P6,可选接口)----
// KnownAttempt 是 pay 合法签发过的一笔尝试的对账维度:期望金额 = base(AmountMinor)+ 渠道尾数
// (尾数封在 provider_ref,由渠道自解,pay 不算)。渠道据此判断一笔到账是否"有主"。
type KnownAttempt struct {
AmountMinor int64
ProviderRef string
}
// OrphanScanRequest 扫描某账户 Since 以来、不匹配任何 Known 的到账。
type OrphanScanRequest struct {
AccountID string
Since time.Time
Known []KnownAttempt
}
// OrphanTransfer 一笔"有钱到账但无主"的转账(付错金额/手动转/超窗迟到旧款)。
type OrphanTransfer struct {
TxID string
AmountMinor int64
Currency string
At time.Time
}
// OrphanScanner 自托管渠道(crypto)可选实现:发现到账但不匹配任何 attempt 的转账。
// 网关侧渠道(alipay/stripe)以对账单核对,不实现此接口。
type OrphanScanner interface {
ScanOrphans(ctx context.Context, req OrphanScanRequest) ([]OrphanTransfer, error)
}
// Registry — 方法名 → Provider(设计 §2 Provider adapter 注册表)。启动期注册,运行期只读。
type Registry struct{ providers map[string]Provider }