Files
wangjia da8bbefe2e merge: P8 订阅/recurring + 拒付 chargeback 并入(订阅生命周期/续费/取消/past_due/chargeback/事件集收口)
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013nMthbVEmQquxBRKb9Fj8u

# Conflicts:
#	internal/model/testdb.go
#	internal/provider/provider.go
#	internal/router/router.go
#	internal/store/order_query_test.go
#	main.go
2026-07-10 19:43:32 +08:00

237 lines
7.7 KiB
Go

// Package webhook delivers pay→business webhooks (v2, with event_type) from an
// outbox, signed with the shared HMAC scheme (双向验签,与 v1 notifyBiz 一致),
// with a background retry sweep as the backstop for lost/failed deliveries.
package webhook
import (
"bytes"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"strconv"
"strings"
"time"
"unicode/utf8"
"github.com/google/uuid"
"github.com/wangjia/pay/config"
"github.com/wangjia/pay/internal/store"
"github.com/wangjia/pay/internal/util"
)
// BizConfigFunc resolves a business system's callback URL + HMAC secret.
type BizConfigFunc func(system string) (config.BizSystemConfig, bool)
// OrderPaidFunc reports whether an order is settled (paid). Delivery gate:
// settle 是"先入队后翻转",崩溃窗口里 outbox 可能存在"未付单"的行——投递前必须
// 门禁,否则会把 payment.succeeded 发给业务方、白给权益。
type OrderPaidFunc func(outTradeNo string) (bool, error)
type Notifier struct {
deliveries *store.WebhookStore
bizConfig BizConfigFunc
orderPaid OrderPaidFunc
client *http.Client
now func() time.Time
maxAttempts int
baseBackoff time.Duration
maxBackoff time.Duration
alert func(d *store.WebhookDeliveryRow, reason string)
}
// Option customizes a Notifier's retry/backoff/alerting behavior (functional options).
type Option func(*Notifier)
// WithClock overrides the time source (tests inject a fake clock to drive backoff windows).
func WithClock(f func() time.Time) Option { return func(n *Notifier) { n.now = f } }
// WithMaxAttempts sets how many failed attempts before a delivery is marked dead.
func WithMaxAttempts(m int) Option { return func(n *Notifier) { n.maxAttempts = m } }
// WithBaseBackoff sets the base duration for exponential backoff (attempt 1).
func WithBaseBackoff(d time.Duration) Option { return func(n *Notifier) { n.baseBackoff = d } }
// WithMaxBackoff caps the exponential backoff duration.
func WithMaxBackoff(d time.Duration) Option { return func(n *Notifier) { n.maxBackoff = d } }
// WithAlerter overrides the dead-letter alert hook (default logs).
func WithAlerter(a func(d *store.WebhookDeliveryRow, reason string)) Option {
return func(n *Notifier) { n.alert = a }
}
// NewNotifier builds a Notifier. The original 3-arg call form keeps compiling
// (opts is variadic); pass Option values to customize clock/backoff/alerting.
func NewNotifier(ws *store.WebhookStore, bizConfig BizConfigFunc, orderPaid OrderPaidFunc, opts ...Option) *Notifier {
n := &Notifier{
deliveries: ws, bizConfig: bizConfig, orderPaid: orderPaid,
client: &http.Client{Timeout: 10 * time.Second},
now: time.Now,
maxAttempts: 12, // ~ 覆盖数小时退避后放弃(见 backoffFor 封顶)
baseBackoff: 30 * time.Second, // 首次失败退避基
maxBackoff: time.Hour, // 单次退避封顶
}
for _, o := range opts {
o(n)
}
if n.alert == nil {
n.alert = func(d *store.WebhookDeliveryRow, reason string) {
log.Printf("[webhook][死信] out_trade_no=%s biz=%s event=%s attempts=%d 放弃投递: %s",
d.OutTradeNo, d.BizSystem, d.EventType, d.Attempts, reason)
}
}
return n
}
// backoffFor 计算第 attempts 次失败后的退避:base·2^(attempts-1),封顶 maxBackoff。
// attempts≥1;移位用循环倍增且封顶,防溢出。
func (n *Notifier) backoffFor(attempts int) time.Duration {
d := n.baseBackoff
for i := 1; i < attempts; i++ {
d *= 2
if d >= n.maxBackoff {
return n.maxBackoff
}
}
if d > n.maxBackoff {
return n.maxBackoff
}
return d
}
// fail 统一失败分流:未达上限→退避重排;达上限→死信+告警。
func (n *Notifier) fail(d *store.WebhookDeliveryRow, msg string) {
attempts := d.Attempts + 1 // 本次即将记的失败次数
if attempts >= n.maxAttempts {
_ = n.deliveries.MarkDead(d.ID, msg)
n.alert(d, msg)
return
}
_ = n.deliveries.ScheduleRetry(d.ID, msg, n.now().Add(n.backoffFor(attempts)))
}
// Enqueue implements gateway.WebhookEnqueuer: serialize the domain payload and
// idempotently persist it to the outbox (delivery happens async). refundID is
// the refund's idempotency dimension (payment events pass "").
func (n *Notifier) Enqueue(outTradeNo, bizSystem, eventType, refundID string, data map[string]any) error {
body, err := json.Marshal(data)
if err != nil {
return fmt.Errorf("webhook.Enqueue marshal: %w", err)
}
return n.deliveries.EnqueueDelivery(outTradeNo, bizSystem, eventType, refundID, string(body))
}
// DeliverPending flushes deliverable rows (未投递、未死信、退避到点); returns how
// many succeeded this pass.
func (n *Notifier) DeliverPending(limit int) (int, error) {
rows, err := n.deliveries.ListDeliverable(n.now(), limit)
if err != nil {
return 0, err
}
ok := 0
for i := range rows {
if n.deliverOne(&rows[i]) {
ok++
}
}
return ok, nil
}
func (n *Notifier) deliverOne(d *store.WebhookDeliveryRow) bool {
// 门禁:订单未付不投递(不计失败,等 settle 翻转后自然放行)。
paid, err := n.orderPaid(d.OutTradeNo)
if err != nil {
log.Printf("[webhook] 投递门禁查单失败 %s: %v", d.OutTradeNo, err)
return false
}
if !paid {
return false
}
cfg, found := n.bizConfig(d.BizSystem)
if !found || cfg.CallbackURL == "" {
n.fail(d, "biz system not configured")
return false
}
if !eventSupported(cfg.SupportedEvents, d.EventType) {
_ = n.deliveries.MarkDelivered(d.ID) // 业务方未订阅该事件:视为已受理,不投递、不重试
return true
}
ts := strconv.FormatInt(time.Now().Unix(), 10)
nonce := uuid.NewString()
sign := util.HMACSign(cfg.Secret, d.BizSystem, ts, nonce, d.Payload)
req, err := http.NewRequest(http.MethodPost, cfg.CallbackURL, bytes.NewReader([]byte(d.Payload)))
if err != nil {
n.fail(d, err.Error())
return false
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Pay-System", d.BizSystem)
req.Header.Set("X-Pay-Event", d.EventType)
req.Header.Set("X-Pay-Timestamp", ts)
req.Header.Set("X-Pay-Nonce", nonce)
req.Header.Set("X-Pay-Sign", sign)
resp, err := n.client.Do(req)
if err != nil {
n.fail(d, err.Error())
return false
}
rb, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
resp.Body.Close()
// 约定:业务方返回 HTTP 200 且响应含 SUCCESS 视为受理(与 v1 一致)。
if resp.StatusCode == http.StatusOK && strings.Contains(strings.ToUpper(string(rb)), "SUCCESS") {
_ = n.deliveries.MarkDelivered(d.ID)
return true
}
n.fail(d, fmt.Sprintf("http %d: %s", resp.StatusCode, truncate(string(rb), 120)))
return false
}
// eventSupported 判断业务方是否声明支持该 event_type(设计 §「接入方显式声明支持事件
// 集」):list 为空 → 仅 payment.succeeded 为真(v1/P2 接入方向后兼容,不会突然收到
// subscription.*/chargeback.received 等新事件把它们搞崩);非空则按声明集精确匹配。
func eventSupported(list []string, ev string) bool {
if len(list) == 0 {
return ev == "payment.succeeded"
}
for _, e := range list {
if e == ev {
return true
}
}
return false
}
// truncate safely truncates a string to n bytes without splitting UTF-8 runes.
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
s = s[:n]
for len(s) > 0 && !utf8.ValidString(s) {
s = s[:len(s)-1]
}
return s
}
// Start runs a background retry sweep (backstop for lost/failed webhooks).
func (n *Notifier) Start(interval time.Duration) {
go func() {
defer func() {
if r := recover(); r != nil {
log.Printf("[webhook] retry sweep panic recovered: %v", r)
}
}()
t := time.NewTicker(interval)
defer t.Stop()
for range t.C {
if _, err := n.DeliverPending(50); err != nil {
log.Printf("[webhook] DeliverPending: %v", err)
}
}
}()
}