From ed5ec5c8948732d7e29cd8f7437b2d25229d5340 Mon Sep 17 00:00:00 2001 From: wangjia <809946525@qq.com> Date: Fri, 10 Jul 2026 16:26:32 +0800 Subject: [PATCH] =?UTF-8?q?feat(v2):=20outbox=20=E9=80=80=E6=AC=BE?= =?UTF-8?q?=E6=84=9F=E7=9F=A5(refund=5Fid=20=E8=BF=9B=E5=94=AF=E4=B8=80?= =?UTF-8?q?=E9=94=AE)+=20=E6=8A=95=E9=80=92=E9=97=A8=E7=A6=81=E6=94=BE?= =?UTF-8?q?=E8=A1=8C=E9=80=80=E6=AC=BE=E6=80=81=20+=20Enqueue=20=E8=B4=AF?= =?UTF-8?q?=E9=80=9A=20refundID?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/gateway/gateway.go | 3 ++- internal/gateway/gateway_test.go | 2 +- internal/gateway/settle.go | 2 +- internal/handler/gateway_test.go | 2 +- internal/model/v2.go | 20 +++++++++++++---- internal/model/v2_test.go | 14 ++++++++++++ internal/model/webhook_delivery.go | 5 +++-- internal/store/webhook.go | 12 +++++----- internal/store/webhook_test.go | 35 +++++++++++++++++++++++++----- internal/webhook/notifier.go | 7 +++--- internal/webhook/notifier_test.go | 6 ++--- main.go | 4 ++-- 12 files changed, 84 insertions(+), 28 deletions(-) diff --git a/internal/gateway/gateway.go b/internal/gateway/gateway.go index 33a0a7a..4d4fc0b 100644 --- a/internal/gateway/gateway.go +++ b/internal/gateway/gateway.go @@ -33,8 +33,9 @@ type ProductResolver interface { } // WebhookEnqueuer receives a domain payload to deliver to the business system. +// refundID 为退款事件的幂等维度(payment 事件传 "")。 type WebhookEnqueuer interface { - Enqueue(outTradeNo, bizSystem, eventType string, data map[string]any) error + Enqueue(outTradeNo, bizSystem, eventType, refundID string, data map[string]any) error } type Gateway struct { diff --git a/internal/gateway/gateway_test.go b/internal/gateway/gateway_test.go index b7e49e0..6b27368 100644 --- a/internal/gateway/gateway_test.go +++ b/internal/gateway/gateway_test.go @@ -38,7 +38,7 @@ type spyEnqueuer struct { failNext bool // 置 true 模拟 outbox 入队失败(settle 崩溃窗口测试用) } -func (s *spyEnqueuer) Enqueue(outTradeNo, bizSystem, eventType string, data map[string]any) error { +func (s *spyEnqueuer) Enqueue(outTradeNo, bizSystem, eventType, refundID string, data map[string]any) error { if s.failNext { s.failNext = false return errors.New("outbox down") diff --git a/internal/gateway/settle.go b/internal/gateway/settle.go index 0e1f46a..2526af5 100644 --- a/internal/gateway/settle.go +++ b/internal/gateway/settle.go @@ -92,7 +92,7 @@ func (g *Gateway) enqueuePaymentSucceeded(att *model.Attempt, paidAt time.Time) "channel": att.Channel, "paid_at": paidAt.Format(time.RFC3339), } - return g.webhook.Enqueue(o.OutTradeNo, o.BizSystem, "payment.succeeded", data) + return g.webhook.Enqueue(o.OutTradeNo, o.BizSystem, "payment.succeeded", "", data) } // HandleCallback runs a channel's raw callback through its Provider.VerifyCallback diff --git a/internal/handler/gateway_test.go b/internal/handler/gateway_test.go index 7b8eec2..9763cbc 100644 --- a/internal/handler/gateway_test.go +++ b/internal/handler/gateway_test.go @@ -24,7 +24,7 @@ import ( type nopEnqueuer struct{} -func (nopEnqueuer) Enqueue(string, string, string, map[string]any) error { return nil } +func (nopEnqueuer) Enqueue(string, string, string, string, map[string]any) error { return nil } type oneResolver struct{} diff --git a/internal/model/v2.go b/internal/model/v2.go index b73afb6..94665b0 100644 --- a/internal/model/v2.go +++ b/internal/model/v2.go @@ -32,12 +32,24 @@ const ( type RefundStatus string const ( - RefundRequested RefundStatus = "requested" - RefundProcessing RefundStatus = "processing" - RefundSucceeded RefundStatus = "succeeded" - RefundFailed RefundStatus = "failed" + RefundRequested RefundStatus = "requested" + RefundProcessing RefundStatus = "processing" + RefundSucceeded RefundStatus = "succeeded" + RefundFailed RefundStatus = "failed" + RefundManualPending RefundStatus = "manual_pending" // crypto 自托管:待运营人工 sweep 退款 ) +// Settled 报告订单是否「已真正收到过钱」(paid 及其后的退款态)。webhook 投递门禁用: +// 只有已结算订单的 outbox 才放行——payment 事件挡住崩溃窗口里的未付单,refund 事件 +// 则因订单已 paid 过而正常放行(不会因订单转入退款态被卡)。 +func (s OrderStatusV2) Settled() bool { + switch s { + case OrderPaidV2, OrderRefundingV2, OrderPartRefundedV2, OrderRefundedV2: + return true + } + return false +} + // ---- 业务订单(购买账本)---- type OrderV2 struct { diff --git a/internal/model/v2_test.go b/internal/model/v2_test.go index f698420..1e2536b 100644 --- a/internal/model/v2_test.go +++ b/internal/model/v2_test.go @@ -33,3 +33,17 @@ func TestV2Migrate(t *testing.T) { t.Fatalf("got %+v", got) } } + +func TestOrderStatusSettled(t *testing.T) { + settled := []model.OrderStatusV2{model.OrderPaidV2, model.OrderRefundingV2, model.OrderPartRefundedV2, model.OrderRefundedV2} + for _, s := range settled { + if !s.Settled() { + t.Fatalf("%s should be settled", s) + } + } + for _, s := range []model.OrderStatusV2{model.OrderCreatedV2, model.OrderPendingV2, model.OrderCanceledV2, model.OrderExpiredV2} { + if s.Settled() { + t.Fatalf("%s should NOT be settled", s) + } + } +} diff --git a/internal/model/webhook_delivery.go b/internal/model/webhook_delivery.go index ba1c1c4..f313bd6 100644 --- a/internal/model/webhook_delivery.go +++ b/internal/model/webhook_delivery.go @@ -1,11 +1,12 @@ package model -// WebhookDelivery 是 pay→业务方 webhook 的 outbox(v2)。unique(out_trade_no,event_type) -// 保证同一订单同一事件只入队一次(幂等);后台 Notifier 扫 Delivered=false 重试兜底。 +// WebhookDelivery 是 pay→业务方 webhook 的 outbox(v2)。unique(out_trade_no,event_type,refund_id) +// 保证同一订单同一事件(同一退款单)只入队一次(幂等);后台 Notifier 扫 Delivered=false 重试兜底。 type WebhookDelivery struct { Base OutTradeNo string `gorm:"size:64;not null;uniqueIndex:uq_delivery" json:"out_trade_no"` EventType string `gorm:"size:32;not null;uniqueIndex:uq_delivery" json:"event_type"` + RefundID string `gorm:"size:64;uniqueIndex:uq_delivery" json:"refund_id,omitempty"` // 退款事件的幂等维度;payment 事件为空 BizSystem string `gorm:"index;size:32" json:"biz_system"` Payload string `gorm:"type:text" json:"payload"` // 已序列化的领域 JSON(含 event_type) Delivered bool `gorm:"index;default:false" json:"delivered"` diff --git a/internal/store/webhook.go b/internal/store/webhook.go index 161db96..f253898 100644 --- a/internal/store/webhook.go +++ b/internal/store/webhook.go @@ -17,14 +17,16 @@ type WebhookStore struct{ db *gorm.DB } func NewWebhookStore(db *gorm.DB) *WebhookStore { return &WebhookStore{db: db} } -// EnqueueDelivery inserts an outbox row; a duplicate (out_trade_no,event_type) -// is a no-op (idempotent enqueue) via ON CONFLICT DO NOTHING. -func (s *WebhookStore) EnqueueDelivery(outTradeNo, bizSystem, eventType, payload string) error { +// EnqueueDelivery inserts an outbox row; a duplicate (out_trade_no,event_type,refund_id) +// is a no-op (idempotent enqueue) via ON CONFLICT DO NOTHING. refundID is the +// refund's idempotency dimension (payment events pass ""). +func (s *WebhookStore) EnqueueDelivery(outTradeNo, bizSystem, eventType, refundID, payload string) error { row := model.WebhookDelivery{ - OutTradeNo: outTradeNo, BizSystem: bizSystem, EventType: eventType, Payload: payload, + OutTradeNo: outTradeNo, BizSystem: bizSystem, EventType: eventType, + RefundID: refundID, Payload: payload, } err := s.db.Clauses(clause.OnConflict{ - Columns: []clause.Column{{Name: "out_trade_no"}, {Name: "event_type"}}, + Columns: []clause.Column{{Name: "out_trade_no"}, {Name: "event_type"}, {Name: "refund_id"}}, DoNothing: true, }).Create(&row).Error if err != nil { diff --git a/internal/store/webhook_test.go b/internal/store/webhook_test.go index 162e3bc..3c9cbef 100644 --- a/internal/store/webhook_test.go +++ b/internal/store/webhook_test.go @@ -11,11 +11,11 @@ import ( func TestWebhookOutboxEnqueueIdempotent(t *testing.T) { ws := store.NewWebhookStore(model.OpenTestDB(t)) - if err := ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", `{"a":1}`); err != nil { + if err := ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", "", `{"a":1}`); err != nil { t.Fatalf("enqueue#1: %v", err) } - // 幂等:同 (out_trade_no,event_type) 再入队不新增行、不报错。 - if err := ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", `{"a":1}`); err != nil { + // 幂等:同 (out_trade_no,event_type,refund_id) 再入队不新增行、不报错。 + if err := ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", "", `{"a":1}`); err != nil { t.Fatalf("enqueue#2: %v", err) } list, _ := ws.ListUndelivered(10) @@ -33,7 +33,7 @@ func TestWebhookOutboxEnqueueIdempotent(t *testing.T) { func TestWebhookMarkFailed(t *testing.T) { ws := store.NewWebhookStore(model.OpenTestDB(t)) - _ = ws.EnqueueDelivery("PAY-2", "jiu", "payment.succeeded", `{}`) + _ = ws.EnqueueDelivery("PAY-2", "jiu", "payment.succeeded", "", `{}`) list, _ := ws.ListUndelivered(10) if err := ws.MarkFailed(list[0].ID, "boom"); err != nil { t.Fatalf("markFailed: %v", err) @@ -46,7 +46,7 @@ func TestWebhookMarkFailed(t *testing.T) { func TestWebhookMarkFailedUTF8Truncation(t *testing.T) { ws := store.NewWebhookStore(model.OpenTestDB(t)) - _ = ws.EnqueueDelivery("PAY-3", "pangolin", "payment.succeeded", `{}`) + _ = ws.EnqueueDelivery("PAY-3", "pangolin", "payment.succeeded", "", `{}`) list, _ := ws.ListUndelivered(10) // Create a long Chinese error message that exceeds 255 bytes @@ -79,3 +79,28 @@ func TestWebhookMarkFailedUTF8Truncation(t *testing.T) { t.Logf("UTF-8 truncated error (%d bytes): %q", len(stored), stored) } + +func TestEnqueueDeliveryRefundIDUnique(t *testing.T) { + ws := store.NewWebhookStore(model.OpenTestDB(t)) + // payment:refund_id="" —— 同单同事件第二次幂等丢弃 + must := func(err error) { + if err != nil { + t.Fatalf("enqueue: %v", err) + } + } + must(ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", "", `{"a":1}`)) + must(ws.EnqueueDelivery("PAY-1", "pangolin", "payment.succeeded", "", `{"a":2}`)) // 幂等 no-op + // refund:同单同事件、不同 refund_id —— 各自成行(部分退多次) + must(ws.EnqueueDelivery("PAY-1", "pangolin", "refund.succeeded", "rf-A", `{"r":"A"}`)) + must(ws.EnqueueDelivery("PAY-1", "pangolin", "refund.succeeded", "rf-B", `{"r":"B"}`)) + must(ws.EnqueueDelivery("PAY-1", "pangolin", "refund.succeeded", "rf-A", `{"r":"A2"}`)) // 同 refund_id 幂等 + + rows, err := ws.ListUndelivered(50) + if err != nil { + t.Fatal(err) + } + // 期望 3 行:1 payment + 2 refund(rf-A / rf-B) + if len(rows) != 3 { + t.Fatalf("undelivered rows = %d want 3: %+v", len(rows), rows) + } +} diff --git a/internal/webhook/notifier.go b/internal/webhook/notifier.go index 67ed598..3b948fa 100644 --- a/internal/webhook/notifier.go +++ b/internal/webhook/notifier.go @@ -43,13 +43,14 @@ func NewNotifier(ws *store.WebhookStore, bizConfig BizConfigFunc, orderPaid Orde } // Enqueue implements gateway.WebhookEnqueuer: serialize the domain payload and -// idempotently persist it to the outbox (delivery happens async). -func (n *Notifier) Enqueue(outTradeNo, bizSystem, eventType string, data map[string]any) error { +// 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, string(body)) + return n.deliveries.EnqueueDelivery(outTradeNo, bizSystem, eventType, refundID, string(body)) } // DeliverPending flushes undelivered rows; returns how many succeeded this pass. diff --git a/internal/webhook/notifier_test.go b/internal/webhook/notifier_test.go index 3a743c7..abcdb39 100644 --- a/internal/webhook/notifier_test.go +++ b/internal/webhook/notifier_test.go @@ -37,7 +37,7 @@ func TestNotifierDeliversSignedEvent(t *testing.T) { n := webhook.NewNotifier(ws, bizCfg, alwaysPaid) // 经 Enqueuer 接口入队(gateway 就是这么调的)。 - err := n.Enqueue("PAY-1", "pangolin", "payment.succeeded", map[string]any{ + err := n.Enqueue("PAY-1", "pangolin", "payment.succeeded", "", map[string]any{ "event_type": "payment.succeeded", "out_trade_no": "PAY-1", "amount_minor": 29990000, "currency": "USDT", }) if err != nil { @@ -84,7 +84,7 @@ func TestNotifierRetriesOnFailure(t *testing.T) { n := webhook.NewNotifier(ws, func(string) (config.BizSystemConfig, bool) { return config.BizSystemConfig{CallbackURL: srv.URL, Secret: "x"}, true }, func(string) (bool, error) { return true, nil }) - _ = n.Enqueue("PAY-3", "pangolin", "payment.succeeded", map[string]any{"event_type": "payment.succeeded"}) + _ = n.Enqueue("PAY-3", "pangolin", "payment.succeeded", "", map[string]any{"event_type": "payment.succeeded"}) if sent, _ := n.DeliverPending(10); sent != 0 { t.Fatalf("失败不应算投递成功, got %d", sent) @@ -115,7 +115,7 @@ func TestNotifierGateSkipsUnpaidOrder(t *testing.T) { n := webhook.NewNotifier(ws, func(string) (config.BizSystemConfig, bool) { return config.BizSystemConfig{CallbackURL: srv.URL, Secret: "x"}, true }, func(string) (bool, error) { return paid, nil }) - _ = n.Enqueue("PAY-4", "pangolin", "payment.succeeded", map[string]any{"event_type": "payment.succeeded"}) + _ = n.Enqueue("PAY-4", "pangolin", "payment.succeeded", "", map[string]any{"event_type": "payment.succeeded"}) if sent, _ := n.DeliverPending(10); sent != 0 || hits != 0 { t.Fatalf("未付单不应投递, sent=%d hits=%d", sent, hits) diff --git a/main.go b/main.go index 104a679..4519b17 100644 --- a/main.go +++ b/main.go @@ -41,11 +41,11 @@ func main() { orderStore := store.NewOrderStore(db) webhookStore := store.NewWebhookStore(db) notifier := webhook.NewNotifier(webhookStore, config.C.BizByName, func(no string) (bool, error) { - o, err := orderStore.GetOrder(no) // 投递门禁:订单已付才发(enqueue-before-flip 不变量) + o, err := orderStore.GetOrder(no) // 投递门禁:订单已结算(paid 及退款态)才放行 if err != nil { return false, err } - return o.Status == model.OrderPaidV2, nil + return o.Status.Settled(), nil }) notifier.Start(60 * time.Second) productResolver := gateway.NewDBProductResolver(db) // 币种由下单渠道结算能力驱动(设计 §4.1)