package gateway import ( "context" "errors" "fmt" "log" "time" "github.com/wangjia/pay/internal/model" "github.com/wangjia/pay/internal/provider" "github.com/wangjia/pay/internal/store" "github.com/wangjia/pay/internal/util" "github.com/wangjia/pay/internal/accounts" ) // webhook event 常量集中在 internal/gateway(Task 7 收敛业务方声明)。 const ( EvtPaymentSucceeded = "payment.succeeded" EvtSubscriptionCreated = "subscription.created" EvtSubscriptionRenewed = "subscription.renewed" EvtSubscriptionPastDue = "subscription.past_due" EvtSubscriptionCanceled = "subscription.canceled" EvtChargebackReceived = "chargeback.received" ) type CreateSubscriptionInput struct { SKU string Method string BizSystem string BizRef string ReturnURL string } type SubscriptionResult struct { SubID string `json:"sub_id"` OrderNo string `json:"order_no"` Session SessionView `json:"session"` } // CreateSubscription 解析套餐权威金额 → 选 stripe 账户(须实现 SubscriptionProvider 且 // SupportsRecurring)→ 落 pending 首购 OrderV2 + Attempt → 返回订阅 Checkout redirect。 // 订阅在首期支付回调时诞生(见 onSubscriptionActivated)。 func (g *Gateway) CreateSubscription(ctx context.Context, in CreateSubscriptionInput) (*SubscriptionResult, error) { prov, err := g.providers.Get(in.Method) if err != nil { return nil, err } caps := prov.Capabilities() subProv, ok := prov.(provider.SubscriptionProvider) if !ok || !caps.SupportsRecurring { return nil, provider.ErrNotSupported } if len(caps.SettleCurrencies) == 0 { return nil, ErrNoSettleCurrency } currency := caps.SettleCurrencies[0] amount, subject, bizCode, err := g.products.Resolve(in.SKU, currency) if err != nil { return nil, err // ErrProductNotFound(含"该币种无价") } outTradeNo := util.NewOutTradeNo("pay") acct, err := g.picker.Pick(in.Method, g.region, accounts.PickHint{OutTradeNo: outTradeNo, AmountMinor: amount}) if err != nil { if errors.Is(err, accounts.ErrNoAccount) { return nil, ErrNoAccount } return nil, err } // SubID 确定性派生自 out_trade_no(与 onSubscriptionActivated 一致,消除双号)。 subID := "SUB-" + outTradeNo sess, err := subProv.CreateSubscriptionCheckout(ctx, provider.CreateRequest{ OutTradeNo: outTradeNo, Subject: subject, AmountMinor: amount, Currency: currency, Account: acct, ReturnURL: in.ReturnURL, Metadata: map[string]string{"pay_sub_id": subID}, }) if err != nil { return nil, fmt.Errorf("gateway.CreateSubscription: %w", err) } order := &model.OrderV2{ OutTradeNo: outTradeNo, BizSystem: in.BizSystem, BizRef: in.BizRef, BizCode: bizCode, Subject: subject, AmountMinor: amount, Currency: currency, Status: model.OrderPendingV2, } if err := g.orders.CreateOrder(order); err != nil { return nil, err } att := &model.Attempt{ OutTradeNo: outTradeNo, Channel: in.Method, AccountID: acct.AccountID, Provider: in.Method, ProviderRef: sess.ProviderRef, RenderType: string(sess.RenderType), AmountMinor: amount, Currency: currency, Status: model.AttemptPending, ExpiresAt: sess.ExpiresAt, } if err := g.orders.CreateAttempt(att); err != nil { return nil, err } return &SubscriptionResult{ SubID: subID, OrderNo: outTradeNo, Session: SessionView{RenderType: string(sess.RenderType), Payload: sess.Payload, ExpiresAt: sess.ExpiresAt}, }, nil } // onSubscriptionActivated 幂等诞生订阅 + 入队 subscription.created。首期支付回调触发。 // created 事件走首购 order 的 out_trade_no + event_type=subscription.created(唯一键天然不撞 payment.succeeded)。 func (g *Gateway) onSubscriptionActivated(ctx context.Context, ev *provider.PaidEvent) error { att, err := g.orders.AttemptByProviderRef(ev.ProviderRef) if err != nil { return nil // 首期会话未落库(不该发生);交由 Settle 侧日志,订阅侧静默 } o, err := g.orders.GetOrder(att.OutTradeNo) if err != nil { return err } if !o.Status.Settled() { // 纵深防御:HandleCallback 的门(res==Processed/Duplicate)已挡住未结算订单, // 这里再核一遍订单状态,防止未来调用点绕过门直接调用本函数。 log.Printf("[subscription] onSubscriptionActivated 订单未结算 out_trade_no=%s status=%s,拒绝激活", o.OutTradeNo, o.Status) return nil } subID := "SUB-" + att.OutTradeNo // 与 CreateSubscription 同式派生 → 重投算出同一 SubID,Create 幂等 created, err := g.subs.Create(&model.Subscription{ SubID: subID, OutTradeNo: o.OutTradeNo, BizSystem: o.BizSystem, BizRef: o.BizRef, BizCode: o.BizCode, Channel: att.Channel, ProviderSubRef: ev.SubscriptionRef, RecurringKind: provider.RecurringKindGatewayScheduled, AmountMinor: o.AmountMinor, Currency: o.Currency, Status: model.SubActive, }) if err != nil { return err } if !created || o.BizSystem == "" { return nil // 已诞生过(重投)/ 独立收款无业务方回调 } return g.webhook.Enqueue(o.OutTradeNo, o.BizSystem, EvtSubscriptionCreated, "", map[string]any{ "event_type": EvtSubscriptionCreated, "out_trade_no": o.OutTradeNo, "sub_id": subID, "biz_system": o.BizSystem, "biz_ref": o.BizRef, "product_biz_code": o.BizCode, "amount_minor": o.AmountMinor, "currency": o.Currency, "channel": att.Channel, "created_at": time.Now().Format(time.RFC3339), }) } // SubscriptionView — GET /api/v2/subscriptions/:sub_id 返回视图(状态 + 续费锚点)。 type SubscriptionView struct { SubID string `json:"sub_id"` Status string `json:"status"` CurrentPeriodEnd *time.Time `json:"current_period_end,omitempty"` CanceledAt *time.Time `json:"canceled_at,omitempty"` } func (g *Gateway) GetSubscription(subID string) (*SubscriptionView, error) { sub, err := g.subs.GetBySubID(subID) if err != nil { return nil, err // store.ErrSubNotFound } return &SubscriptionView{ SubID: sub.SubID, Status: string(sub.Status), CurrentPeriodEnd: sub.CurrentPeriodEnd, CanceledAt: sub.CanceledAt, }, nil } // CancelSubscription 主动取消(POST /api/v2/subscriptions/:sub_id/cancel):查订阅 → 已终态则 // 幂等 no-op(不再打渠道,避免重复取消命中渠道 400)→ 渠道侧取消 → 本地翻 canceled + 入队 // subscription.canceled。Stripe 随后异步发 customer.subscription.deleted,入站处理器 // (settleSubscriptionCanceled)再次调用同一 finalizeCanceled,两路收敛同一终态、天然幂等。 func (g *Gateway) CancelSubscription(ctx context.Context, subID string) error { sub, err := g.subs.GetBySubID(subID) if err != nil { return err // store.ErrSubNotFound } if sub.Status == model.SubCanceled { return nil // 已取消(重复调用/webhook 已先到):幂等 no-op,不再打渠道 } prov, err := g.providers.Get(sub.Channel) if err != nil { return err } if sp, ok := prov.(provider.SubscriptionProvider); ok { if err := sp.CancelSubscription(ctx, sub.ProviderSubRef); err != nil { return err } } return g.finalizeCanceled(sub) } // settleSubscriptionCanceled 处理入站 customer.subscription.deleted:反查订阅 → finalizeCanceled。 // 未知 provider_sub_ref(订阅未诞生/已清理)忽略,不报错(与 settleRenewal 同惯例)。 func (g *Gateway) settleSubscriptionCanceled(ctx context.Context, method string, ev *provider.PaidEvent) (SettleResult, error) { sub, err := g.subs.GetByProviderRef(method, ev.SubscriptionRef) if err != nil { if errors.Is(err, store.ErrSubNotFound) { return SettleIgnored, nil } return SettleFailed, err } if err := g.finalizeCanceled(sub); err != nil { return SettleFailed, err } return SettleProcessed, nil } // finalizeCanceled 是取消的唯一落地点(主动取消 API + 入站 webhook 两路收敛于此): // MarkCanceled 幂等翻转 canceled;仅首次真正翻转(flipped)且有业务方(BizSystem 非空) // 才入队 subscription.canceled——重投/独立收款不重复发。 // // outbox 键取真实首购单号 sub.OutTradeNo(不用合成后缀键):该单已 paid,Notifier 投递门禁 // (orderPaid)天然放行;event_type=subscription.canceled 与该单已有的 payment.succeeded / // subscription.created 行不同 event_type,唯一键 (out_trade_no,event_type,refund_id) 不撞。 // canceled 每订阅只发生一次,不存在"同订阅多次 canceled 被吞"的问题(不同于 past_due)。 func (g *Gateway) finalizeCanceled(sub *model.Subscription) error { flipped, err := g.subs.MarkCanceled(sub.SubID) if err != nil { return err } if !flipped || sub.BizSystem == "" { return nil // 已取消(重投)/ 独立收款无业务方回调 } return g.webhook.Enqueue(sub.OutTradeNo, sub.BizSystem, EvtSubscriptionCanceled, "", map[string]any{ "event_type": EvtSubscriptionCanceled, "out_trade_no": sub.OutTradeNo, "sub_id": sub.SubID, "biz_system": sub.BizSystem, "biz_ref": sub.BizRef, "product_biz_code": sub.BizCode, "canceled_at": time.Now().Format(time.RFC3339), }) } // markSubscriptionPastDue 处理入站 invoice.payment_failed:active→past_due(MarkPastDue 只在 // 当前 active 时翻转,已 past_due/已 canceled 不动),入队 subscription.past_due 供业务方提醒 // 用户换卡。下期 invoice.paid 成功由 settleRenewal 的 Activate 自动恢复 active。 // // outbox 键取真实首购单号 sub.OutTradeNo(理由同 finalizeCanceled):Task5 最小可行接受 // "同订阅多次 past_due 只报首次"(第二次失败的行撞 (out_trade_no,event_type) 唯一键、 // EnqueueDelivery 幂等 no-op),Task 7 视需要再引入合成键 + Notifier 门禁旁路。 func (g *Gateway) markSubscriptionPastDue(ctx context.Context, method string, ev *provider.PaidEvent) (SettleResult, error) { sub, err := g.subs.GetByProviderRef(method, ev.SubscriptionRef) if err != nil { if errors.Is(err, store.ErrSubNotFound) { return SettleIgnored, nil } return SettleFailed, err } flipped, err := g.subs.MarkPastDue(sub.Channel, sub.ProviderSubRef) if err != nil { return SettleFailed, err } if !flipped || sub.BizSystem == "" { return SettleDuplicate, nil } if err := g.webhook.Enqueue(sub.OutTradeNo, sub.BizSystem, EvtSubscriptionPastDue, "", map[string]any{ "event_type": EvtSubscriptionPastDue, "out_trade_no": sub.OutTradeNo, "sub_id": sub.SubID, "biz_system": sub.BizSystem, "biz_ref": sub.BizRef, "product_biz_code": sub.BizCode, "failed_at": time.Now().Format(time.RFC3339), }); err != nil { return SettleFailed, err } return SettleProcessed, nil }