diff --git a/client/agpay_client.go b/client/agpay_client.go new file mode 100644 index 0000000..706f025 --- /dev/null +++ b/client/agpay_client.go @@ -0,0 +1,66 @@ +package client + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "strings" + "time" +) + +type AGPayClient struct { + baseURL string + httpClient *http.Client +} + +type AGPayCreateRequest struct { + TotalFeeFen int64 `json:"totalFee"` + MerchantOrder string `json:"tradeNo"` + NotifyURL string `json:"notifyUrl"` +} + +type AGPayRequestRejectedError struct { + StatusCode int +} + +func (e *AGPayRequestRejectedError) Error() string { + return fmt.Sprintf("AGPay rejected payment request with HTTP %d", e.StatusCode) +} + +func NewAGPayClient(baseURL string) (*AGPayClient, error) { + if baseURL == "" { + return nil, errors.New("AGPay base URL is required") + } + return &AGPayClient{baseURL: strings.TrimRight(baseURL, "/"), httpClient: &http.Client{Timeout: 15 * time.Second}}, nil +} + +func (c *AGPayClient) CreateWechatQRCode(req AGPayCreateRequest) ([]byte, error) { + if req.TotalFeeFen <= 0 || req.MerchantOrder == "" || req.NotifyURL == "" { + return nil, errors.New("invalid AGPay payment request") + } + body, err := json.Marshal(req) + if err != nil { + return nil, err + } + httpReq, err := http.NewRequest(http.MethodPost, c.baseURL+"/pay/wx/pay", bytes.NewReader(body)) + if err != nil { + return nil, err + } + httpReq.Header.Set("Content-Type", "application/json") + resp, err := c.httpClient.Do(httpReq) + if err != nil { + return nil, fmt.Errorf("call AGPay: %w", err) + } + defer resp.Body.Close() + response, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if err != nil { + return nil, err + } + if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices { + return nil, &AGPayRequestRejectedError{StatusCode: resp.StatusCode} + } + return response, nil +} diff --git a/client/simboss_client.go b/client/simboss_client.go new file mode 100644 index 0000000..d589c30 --- /dev/null +++ b/client/simboss_client.go @@ -0,0 +1,156 @@ +package client + +import ( + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "sort" + "strconv" + "strings" + "time" +) + +const defaultSimbossAPIBase = "https://api.simboss.com" + +type SimbossClient struct { + appID string + secret string + apiBase string + httpClient *http.Client +} + +type SimbossDeviceDetail struct { + Carrier string `json:"carrier"` + Status string `json:"status"` + DeviceStatus string `json:"deviceStatus"` + ExpireDate string `json:"expireDate"` + RatePlanID int64 `json:"ratePlanId"` + RatePlanName string `json:"iratePlanName"` + DataUsage float64 `json:"dataUsage"` + TotalDataVolume float64 `json:"totalDataVolume"` + RatePlanExpirationDate string `json:"ratePlanExpirationDate"` +} + +type SimbossRatePlan struct { + RatePlanID int64 `json:"ratePlanId"` + Name string `json:"name"` + Description string `json:"description"` + DataVolume float64 `json:"dataVolume"` + TimeLength int `json:"timeLength"` + TimeUnit string `json:"timeUnit"` + MaxRechargePeriod int `json:"maxRechargePeriod"` +} + +type simbossResponse struct { + Code string `json:"code"` + Message string `json:"message"` + Success bool `json:"success"` + Data json.RawMessage `json:"data"` +} + +func NewSimbossClient(appID, secret, apiBase string) (*SimbossClient, error) { + if appID == "" || secret == "" { + return nil, errors.New("SIMBOSS credentials are required") + } + if apiBase == "" { + apiBase = defaultSimbossAPIBase + } + return &SimbossClient{ + appID: appID, secret: secret, apiBase: strings.TrimRight(apiBase, "/"), + httpClient: &http.Client{Timeout: 30 * time.Second}, + }, nil +} + +func (c *SimbossClient) GetDeviceDetail(iccid string) (*SimbossDeviceDetail, error) { + var detail SimbossDeviceDetail + if err := c.post("/2.0/device/detail", map[string]string{"iccid": iccid}, &detail); err != nil { + return nil, err + } + return &detail, nil +} + +func (c *SimbossClient) GetRatePlans(iccid string) ([]SimbossRatePlan, error) { + var plans []SimbossRatePlan + if err := c.post("/2.0/device/rateplans", map[string]string{"iccid": iccid}, &plans); err != nil { + return nil, err + } + return plans, nil +} + +func (c *SimbossClient) Recharge(iccid string, ratePlanID int64, months int, externalOrder string) (string, error) { + if ratePlanID <= 0 || months <= 0 || externalOrder == "" { + return "", errors.New("invalid SIMBOSS recharge request") + } + var sequence string + err := c.post("/2.0/device/recharge", map[string]string{ + "iccid": iccid, + "ratePlanId": strconv.FormatInt(ratePlanID, 10), + "month": strconv.Itoa(months), + "externalOrder": externalOrder, + }, &sequence) + return sequence, err +} + +func (c *SimbossClient) post(path string, params map[string]string, target any) error { + params["appid"] = c.appID + params["timestamp"] = strconv.FormatInt(time.Now().UnixMilli(), 10) + params["sign"] = c.sign(params) + + form := url.Values{} + for key, value := range params { + form.Set(key, value) + } + req, err := http.NewRequest(http.MethodPost, c.apiBase+path, strings.NewReader(form.Encode())) + if err != nil { + return fmt.Errorf("create SIMBOSS request: %w", err) + } + req.Header.Set("Content-Type", "application/x-www-form-urlencoded;charset=utf-8") + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("call SIMBOSS: %w", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if err != nil { + return fmt.Errorf("read SIMBOSS response: %w", err) + } + if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices { + return fmt.Errorf("SIMBOSS returned HTTP %d", resp.StatusCode) + } + var result simbossResponse + if err := json.Unmarshal(body, &result); err != nil { + return fmt.Errorf("decode SIMBOSS response: %w", err) + } + if result.Code != "0" && !result.Success { + return fmt.Errorf("SIMBOSS rejected request: code=%s", result.Code) + } + if err := json.Unmarshal(result.Data, target); err != nil { + return fmt.Errorf("decode SIMBOSS data: %w", err) + } + return nil +} + +func (c *SimbossClient) sign(params map[string]string) string { + keys := make([]string, 0, len(params)) + for key := range params { + keys = append(keys, key) + } + sort.Strings(keys) + var builder strings.Builder + for i, key := range keys { + if i > 0 { + builder.WriteByte('&') + } + builder.WriteString(key) + builder.WriteByte('=') + builder.WriteString(params[key]) + } + builder.WriteString(c.secret) + sum := sha256.Sum256([]byte(builder.String())) + return hex.EncodeToString(sum[:]) +} diff --git a/common/busi_error.go b/common/busi_error.go index 2e2659f..58855a7 100644 --- a/common/busi_error.go +++ b/common/busi_error.go @@ -70,6 +70,15 @@ const ( CarrierUnavailable = 58005 PlatformTrafficExhausted = 58006 TrafficLedgerConflict = 58007 + + // 支付 59xxx + PaymentDisabled = 59001 + PaymentOrderConflict = 59002 + PaymentAlreadyPaid = 59003 + PaymentTransactionNotFound = 59004 + PendingOrderExists = 59005 + OrderClosed = 59006 + OrderNotCancellable = 59007 ) // 预定义的业务错误 @@ -116,13 +125,20 @@ var ( ErrVideoNotFound = &BusiError{Code: VideoNotFound, Msg: "视频不存在"} ErrVideoUploading = &BusiError{Code: VideoUploading, Msg: "视频上传中"} - ErrTrafficNotEnough = &BusiError{Code: TrafficNotEnough, Msg: "流量不足,请充值"} - ErrOrderPaid = &BusiError{Code: OrderPaid, Msg: "该订单已支付"} - ErrSimCardNotFound = &BusiError{Code: SimCardNotFound, Msg: "SIM卡不存在"} - ErrOrderNotFound = &BusiError{Code: OrderNotFound, Msg: "订单不存在"} - ErrCarrierUnavailable = &BusiError{Code: CarrierUnavailable, Msg: "运营商通道未接入"} - ErrPlatformTrafficExhausted = &BusiError{Code: PlatformTrafficExhausted, Msg: "资源异常,请联系管理员"} - ErrTrafficLedgerConflict = &BusiError{Code: TrafficLedgerConflict, Msg: "账务流水幂等冲突"} + ErrTrafficNotEnough = &BusiError{Code: TrafficNotEnough, Msg: "流量不足,请充值"} + ErrOrderPaid = &BusiError{Code: OrderPaid, Msg: "该订单已支付"} + ErrSimCardNotFound = &BusiError{Code: SimCardNotFound, Msg: "SIM卡不存在"} + ErrOrderNotFound = &BusiError{Code: OrderNotFound, Msg: "订单不存在"} + ErrCarrierUnavailable = &BusiError{Code: CarrierUnavailable, Msg: "运营商通道未接入"} + ErrPlatformTrafficExhausted = &BusiError{Code: PlatformTrafficExhausted, Msg: "资源异常,请联系管理员"} + ErrTrafficLedgerConflict = &BusiError{Code: TrafficLedgerConflict, Msg: "账务流水幂等冲突"} + ErrPaymentDisabled = &BusiError{Code: PaymentDisabled, Msg: "支付服务未配置"} + ErrPaymentOrderConflict = &BusiError{Code: PaymentOrderConflict, Msg: "支付订单幂等冲突"} + ErrPaymentAlreadyPaid = &BusiError{Code: PaymentAlreadyPaid, Msg: "订单已支付"} + ErrPaymentTransactionNotFound = &BusiError{Code: PaymentTransactionNotFound, Msg: "支付交易不存在"} + ErrPendingOrderExists = &BusiError{Code: PendingOrderExists, Msg: "存在未付款的订单,请前往订单页面付款或取消该订单后再创建新订单"} + ErrOrderClosed = &BusiError{Code: OrderClosed, Msg: "订单已关闭"} + ErrOrderNotCancellable = &BusiError{Code: OrderNotCancellable, Msg: "订单当前状态不允许取消"} ) // NewBusiError 创建自定义业务错误 diff --git a/common/config.go b/common/config.go index 6c4f747..74de4e8 100644 --- a/common/config.go +++ b/common/config.go @@ -17,6 +17,8 @@ type AppConfig struct { Log Log `mapstructure:"log"` OSS OSS `mapstructure:"oss"` Live Live `mapstructure:"live"` + Simboss Simboss `mapstructure:"simboss"` + AGPay AGPay `mapstructure:"agpay"` Heartbeat Heartbeat `mapstructure:"heartbeat"` } @@ -89,6 +91,16 @@ type Live struct { AllowRealCloud bool `mapstructure:"allow-real-cloud"` } +type Simboss struct { + AppID string `mapstructure:"app-id"` + Secret string `mapstructure:"secret"` + APIBase string `mapstructure:"api-base"` +} + +type AGPay struct { + BaseURL string `mapstructure:"base-url"` + NotifyURL string `mapstructure:"notify-url"` +} type Heartbeat struct { TimeoutSeconds int `mapstructure:"timeout-seconds"` ScanSeconds int `mapstructure:"scan-seconds"` diff --git a/config-prod.yaml b/config-prod.yaml index 3f495ad..f8fa885 100644 --- a/config-prod.yaml +++ b/config-prod.yaml @@ -47,6 +47,13 @@ live: record-policy: disabled callback-auth-token: "db570a886e3d65c3f6c55d648ddb41506a658f6be4892719345056602fc2e7ac" allow-real-cloud: true +simboss: + app-id: "1024201105934" + secret: "1fb8296b2e1a0b25010912224c48173b" + api-base: "https://api.simboss.com" +agpay: + base-url: "http://192.168.1.195:8888" + notify-url: "http://192.168.0.6:9913/pay/wx/callback" heartbeat: timeout-seconds: 30 scan-seconds: 5 diff --git a/config.yaml b/config.yaml index 68ec9b9..6ebc275 100644 --- a/config.yaml +++ b/config.yaml @@ -47,6 +47,14 @@ live: billing-interval-seconds: 60 record-policy: disabled allow-real-cloud: false +simboss: + app-id: "1024201105934" + secret: "1fb8296b2e1a0b25010912224c48173b" + api-base: "https://api.simboss.com" +agpay: + base-url: "http://192.168.1.195:8888" + notify-url: "http://192.168.0.6:8080/pay/wx/callback" + heartbeat: timeout-seconds: 30 scan-seconds: 5 diff --git a/handler/account_handler.go b/handler/account_handler.go index 49e8c51..22e68c0 100644 --- a/handler/account_handler.go +++ b/handler/account_handler.go @@ -76,6 +76,15 @@ func GetOrderPage(c *gin.Context) { } common.OKWithData(c, data) } +func ListTrafficPackages(c *gin.Context) { + data, e := service.DefaultAccountService.ListTrafficPackages() + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + func CreateTrafficOrder(c *gin.Context) { var req vo.TrafficOrderCreateReq if err := c.ShouldBindJSON(&req); err != nil { @@ -89,18 +98,17 @@ func CreateTrafficOrder(c *gin.Context) { } common.OKWithData(c, data) } -func PayTrafficOrder(c *gin.Context) { +func CancelTrafficOrder(c *gin.Context) { id, e := parseID(c) if e != nil { common.FailWithBusiError(c, e) return } - data, e := service.DefaultAccountService.PayTrafficOrder(id, common.GetUserId(c), c.GetHeader("Idempotency-Key")) - if e != nil { + if e := service.DefaultAccountService.CancelTrafficOrder(common.GetUserId(c), id); e != nil { common.FailWithBusiError(c, e) return } - common.OKWithData(c, data) + common.OK(c) } func ListSimCards(c *gin.Context) { var req vo.SimCardPageReq @@ -124,6 +132,74 @@ func GetSimRechargeLogPage(c *gin.Context) { } common.OKWithData(c, data) } +func GetSimCard(c *gin.Context) { + id, e := parseID(c) + if e != nil { + common.FailWithBusiError(c, e) + return + } + refresh := c.Query("refresh") == "true" + data, e := service.DefaultSimRechargeService.GetCard(common.GetUserId(c), id, refresh) + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + +func ListSimPackages(c *gin.Context) { + data, e := service.DefaultSimRechargeService.ListPackages() + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + +func CreateSimRechargeOrder(c *gin.Context) { + id, e := parseID(c) + if e != nil { + common.FailWithBusiError(c, e) + return + } + var req vo.SimRechargeReq + if err := c.ShouldBindJSON(&req); err != nil { + common.FailWithBindError(c, common.ErrParam, err) + return + } + data, e := service.DefaultAccountService.RechargeSimCard(common.GetUserId(c), id, &req) + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + +func CancelSimRechargeOrder(c *gin.Context) { + id, e := parseID(c) + if e != nil { + common.FailWithBusiError(c, e) + return + } + if e := service.DefaultAccountService.CancelSimRechargeOrder(common.GetUserId(c), id); e != nil { + common.FailWithBusiError(c, e) + return + } + common.OK(c) +} +func GetSimRechargeOrder(c *gin.Context) { + id, e := parseID(c) + if e != nil { + common.FailWithBusiError(c, e) + return + } + data, e := service.DefaultSimRechargeService.GetOrder(common.GetUserId(c), id) + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} func RechargeSimCard(c *gin.Context) { id, e := parseID(c) if e != nil { diff --git a/handler/payment_handler.go b/handler/payment_handler.go new file mode 100644 index 0000000..dc8ce3e --- /dev/null +++ b/handler/payment_handler.go @@ -0,0 +1,79 @@ +package handler + +import ( + "bytes" + "encoding/json" + "io" + "mime" + "net/http" + + "github.com/gin-gonic/gin" + "github.com/gin-gonic/gin/binding" + + "laic-backend/common" + "laic-backend/model" + "laic-backend/service" + "laic-backend/vo" +) + +func CreateWechatPayment(c *gin.Context) { + var req vo.PaymentCreateReq + if err := c.ShouldBindJSON(&req); err != nil { + common.FailWithBindError(c, common.ErrParam, err) + return + } + data, e := service.DefaultPaymentService.CreateWechatPayment(common.GetUserId(c), &req) + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + +func GetPaymentTransaction(c *gin.Context) { + id, e := parseID(c) + if e != nil { + common.FailWithBusiError(c, e) + return + } + data, e := service.DefaultPaymentService.GetTransaction(common.GetUserId(c), id) + if e != nil { + common.FailWithBusiError(c, e) + return + } + common.OKWithData(c, data) +} + +func AGPayCallback(c *gin.Context) { + contentType, _, err := mime.ParseMediaType(c.GetHeader("Content-Type")) + if err != nil || contentType != "application/json" { + c.Status(http.StatusBadRequest) + return + } + c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, 1<<10) + body, err := c.GetRawData() + if err != nil { + c.Status(http.StatusBadRequest) + return + } + var callback model.AGPayCallback + decoder := json.NewDecoder(bytes.NewReader(body)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&callback); err != nil || decoder.Decode(&struct{}{}) != io.EOF { + c.Status(http.StatusBadRequest) + return + } + if err := binding.Validator.ValidateStruct(&callback); err != nil { + c.Status(http.StatusBadRequest) + return + } + if err := service.DefaultPaymentService.ConfirmAGPayCallback(callback.TradeNo, body); err != nil { + if err == common.ErrInternal { + c.Status(http.StatusInternalServerError) + return + } + c.Status(http.StatusBadRequest) + return + } + c.JSON(http.StatusOK, gin.H{"code": "SUCCESS"}) +} diff --git a/main.go b/main.go index 7840cdb..fa2a2b1 100644 --- a/main.go +++ b/main.go @@ -29,6 +29,14 @@ func main() { logger.ERROR("直播提供商初始化失败", err) panic(err) } + if err := service.InitSimboss(conf.Simboss); err != nil { + logger.ERROR("SIMBOSS 初始化失败", err) + panic(err) + } + if err := service.InitAGPay(conf.AGPay); err != nil { + logger.ERROR("AGPay 初始化失败", err) + panic(err) + } setLogLevel(conf.Log.Level) token.Init(conf.JWT.Secret, conf.JWT.AccessExpireH, conf.JWT.RefreshExpireH) diff --git a/model/billing.go b/model/billing.go index de78634..d7c624f 100644 --- a/model/billing.go +++ b/model/billing.go @@ -12,13 +12,20 @@ type TrafficOrder struct { AmountBytes int64 `gorm:"column:amount_bytes;type:BIGINT;not null" json:"amountBytes"` UnitPrice float64 `gorm:"column:unit_price;type:DECIMAL(6,2)" json:"unitPrice"` TotalPrice float64 `gorm:"column:total_price;type:DECIMAL(10,2)" json:"totalPrice"` + TotalFeeFen int64 `gorm:"column:total_fee_fen;type:BIGINT;not null" json:"totalFeeFen"` PayStatus string `gorm:"column:pay_status;type:VARCHAR(16);default:unpaid" json:"payStatus"` + PaymentProvider string `gorm:"column:payment_provider;type:VARCHAR(32)" json:"paymentProvider"` + ProviderTradeNo *string `gorm:"column:provider_trade_no;type:VARCHAR(128)" json:"-"` CreditStatus string `gorm:"column:credit_status;type:VARCHAR(16);default:pending" json:"creditStatus"` CreditLedgerID int64 `gorm:"column:credit_ledger_id;type:BIGINT" json:"creditLedgerId"` + CreditAttemptCount int `gorm:"column:credit_attempt_count;type:INT;not null" json:"-"` + NextCreditAt *time.Time `gorm:"column:next_credit_at" json:"-"` PaymentIdempotencyKey string `gorm:"column:payment_idempotency_key;type:VARCHAR(128);uniqueIndex:uk_traffic_order_payment_idem" json:"-"` PaidBy int64 `gorm:"column:paid_by;type:BIGINT" json:"paidBy"` CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` PaidAt *time.Time `gorm:"column:paid_at" json:"paidAt"` + ClosedAt *time.Time `gorm:"column:closed_at" json:"closedAt"` + UpdatedAt time.Time `gorm:"column:updated_at" json:"updatedAt"` } func (TrafficOrder) TableName() string { diff --git a/model/payment.go b/model/payment.go new file mode 100644 index 0000000..011ecd3 --- /dev/null +++ b/model/payment.go @@ -0,0 +1,39 @@ +package model + +import "time" + +type PaymentTransaction struct { + ID int64 `gorm:"primaryKey;column:id;type:BIGINT;not null" json:"id"` + UserID int64 `gorm:"column:user_id;type:BIGINT;not null" json:"userId"` + BusinessType string `gorm:"column:business_type;type:VARCHAR(32);not null" json:"businessType"` + BusinessOrderID int64 `gorm:"column:business_order_id;type:BIGINT;not null" json:"businessOrderId"` + Provider string `gorm:"column:provider;type:VARCHAR(32);not null" json:"provider"` + MerchantOrderNo string `gorm:"column:merchant_order_no;type:VARCHAR(128);not null" json:"merchantOrderNo"` + ProviderTradeNo *string `gorm:"column:provider_trade_no;type:VARCHAR(128)" json:"-"` + ProviderPayload string `gorm:"column:provider_payload;type:TEXT" json:"-"` + RequestedAt *time.Time `gorm:"column:requested_at" json:"-"` + TotalFeeFen int64 `gorm:"column:total_fee_fen;type:BIGINT;not null" json:"totalFeeFen"` + Status string `gorm:"column:status;type:VARCHAR(16);not null;default:unpaid" json:"status"` + CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` + PaidAt *time.Time `gorm:"column:paid_at" json:"paidAt"` + UpdatedAt time.Time `gorm:"column:updated_at" json:"updatedAt"` +} + +func (PaymentTransaction) TableName() string { return "payment_transaction" } + +type PaymentCallbackEvent struct { + ID int64 `gorm:"primaryKey;column:id;type:BIGINT;not null" json:"id"` + Provider string `gorm:"column:provider;type:VARCHAR(32);not null" json:"provider"` + EventKey string `gorm:"column:event_key;type:VARCHAR(128);not null" json:"eventKey"` + PayloadHash string `gorm:"column:payload_hash;type:VARCHAR(64);not null" json:"payloadHash"` + Verified bool `gorm:"column:verified;type:TINYINT(1);not null" json:"verified"` + ProcessedAt *time.Time `gorm:"column:processed_at" json:"processedAt"` + ResultCode string `gorm:"column:result_code;type:VARCHAR(64)" json:"resultCode"` + CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` +} + +type AGPayCallback struct { + TradeNo string `json:"tradeNo" binding:"required,max=128"` +} + +func (PaymentCallbackEvent) TableName() string { return "payment_callback_event" } diff --git a/model/sim_recharge.go b/model/sim_recharge.go new file mode 100644 index 0000000..5d927c1 --- /dev/null +++ b/model/sim_recharge.go @@ -0,0 +1,52 @@ +package model + +import "time" + +type SimPackage struct { + ID int64 `gorm:"primaryKey;column:id;type:BIGINT;not null" json:"id"` + Code string `gorm:"column:code;type:VARCHAR(64);not null" json:"code"` + Name string `gorm:"column:name;type:VARCHAR(128);not null" json:"name"` + Carrier string `gorm:"column:carrier;type:VARCHAR(16);not null" json:"carrier"` + ProviderRatePlanID int64 `gorm:"column:provider_rate_plan_id;type:BIGINT;not null" json:"-"` + AmountGb int `gorm:"column:amount_gb;type:INT;not null" json:"amountGb"` + ValidityMonths int `gorm:"column:validity_months;type:INT;not null" json:"validityMonths"` + CostFeeFen int64 `gorm:"column:cost_fee_fen;type:BIGINT;not null" json:"-"` + SaleFeeFen int64 `gorm:"column:sale_fee_fen;type:BIGINT;not null" json:"saleFeeFen"` + Status string `gorm:"column:status;type:VARCHAR(16);not null;default:active" json:"status"` + CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` + UpdatedAt time.Time `gorm:"column:updated_at" json:"updatedAt"` +} + +func (SimPackage) TableName() string { return "sim_package" } + +type SimRechargeOrder struct { + ID int64 `gorm:"primaryKey;column:id;type:BIGINT;not null" json:"id"` + UserID int64 `gorm:"column:user_id;type:BIGINT;not null" json:"userId"` + SimCardID int64 `gorm:"column:sim_card_id;type:BIGINT;not null" json:"simCardId"` + PackageID int64 `gorm:"column:package_id;type:BIGINT;not null" json:"packageId"` + PackageCode string `gorm:"column:package_code;type:VARCHAR(64);not null" json:"packageCode"` + ProviderRatePlanID int64 `gorm:"column:provider_rate_plan_id;type:BIGINT;not null" json:"-"` + IccidSnapshot string `gorm:"column:iccid_snapshot;type:VARCHAR(32);not null" json:"-"` + AmountGb int `gorm:"column:amount_gb;type:INT;not null" json:"amountGb"` + ValidityMonths int `gorm:"column:validity_months;type:INT;not null" json:"validityMonths"` + TotalFeeFen int64 `gorm:"column:total_fee_fen;type:BIGINT;not null" json:"totalFeeFen"` + PaymentStatus string `gorm:"column:payment_status;type:VARCHAR(16);not null;default:unpaid" json:"paymentStatus"` + FulfillmentStatus string `gorm:"column:fulfillment_status;type:VARCHAR(16);not null;default:pending" json:"fulfillmentStatus"` + FulfillmentLeaseToken string `gorm:"column:fulfillment_lease_token;type:VARCHAR(64)" json:"-"` + FulfillmentLeaseUntil *time.Time `gorm:"column:fulfillment_lease_until" json:"-"` + ExternalOrderNo string `gorm:"column:external_order_no;type:VARCHAR(128);not null;uniqueIndex:uk_sim_recharge_external" json:"-"` + ProviderOrderNo string `gorm:"column:provider_order_no;type:VARCHAR(128)" json:"-"` + AttemptCount int `gorm:"column:attempt_count;type:INT;not null;default:0" json:"-"` + NextAttemptAt *time.Time `gorm:"column:next_attempt_at" json:"-"` + LastErrorCode string `gorm:"column:last_error_code;type:VARCHAR(64)" json:"-"` + LastErrorMessage string `gorm:"column:last_error_message;type:VARCHAR(256)" json:"-"` + FulfillmentStartedAt *time.Time `gorm:"column:fulfillment_started_at" json:"-"` + FulfillmentMessage string `gorm:"column:fulfillment_message;type:VARCHAR(128)" json:"fulfillmentMessage"` + CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` + PaidAt *time.Time `gorm:"column:paid_at" json:"paidAt"` + ClosedAt *time.Time `gorm:"column:closed_at" json:"closedAt"` + FulfilledAt *time.Time `gorm:"column:fulfilled_at" json:"fulfilledAt"` + UpdatedAt time.Time `gorm:"column:updated_at" json:"updatedAt"` +} + +func (SimRechargeOrder) TableName() string { return "sim_recharge_order" } diff --git a/route/route.go b/route/route.go index 26ad21f..eaf2b2f 100644 --- a/route/route.go +++ b/route/route.go @@ -24,6 +24,7 @@ func InitRouter(port int32) { }) // 阿里云直播回调使用独立令牌鉴权,不能进入用户 JWT 或操作日志中间件。 + engine.POST("/pay/wx/callback", handler.AGPayCallback) engine.POST("/internal/callbacks/aliyun/live", handler.HandleAliyunLiveCallback) // v1 业务路由 @@ -51,13 +52,21 @@ func InitRouter(port int32) { account.GET("/profile", handler.GetProfile) account.PUT("/profile", handler.UpdateProfile) account.GET("/traffic/balance", handler.GetTrafficBalance) + account.GET("/traffic-packages", handler.ListTrafficPackages) account.GET("/resource-summary", handler.GetResourceSummary) account.GET("/downloads", handler.GetDownloadPage) account.GET("/traffic/usage", handler.GetUsagePage) account.GET("/traffic/orders", handler.GetOrderPage) account.POST("/traffic/orders", handler.CreateTrafficOrder) - account.POST("/traffic/orders/:id/pay", middleware.AdminMiddleware(), handler.PayTrafficOrder) + account.DELETE("/traffic/orders/:id", handler.CancelTrafficOrder) + account.POST("/payments", handler.CreateWechatPayment) + account.GET("/payments/:id", handler.GetPaymentTransaction) account.GET("/sim-cards", handler.ListSimCards) + account.GET("/sim-cards/:id", handler.GetSimCard) + account.GET("/sim-packages", handler.ListSimPackages) + account.POST("/sim-cards/:id/recharge-orders", handler.CreateSimRechargeOrder) + account.GET("/sim-recharge-orders/:id", handler.GetSimRechargeOrder) + account.DELETE("/sim-recharge-orders/:id", handler.CancelSimRechargeOrder) account.GET("/sim-recharge-logs", handler.GetSimRechargeLogPage) account.POST("/security/change-password", handler.ChangePassword) account.GET("/security/sessions", handler.ListSessions) diff --git a/service/account_service.go b/service/account_service.go index 9e242fd..2e2f1fe 100644 --- a/service/account_service.go +++ b/service/account_service.go @@ -79,17 +79,24 @@ func (s *AccountService) GetOrderPage(userID int64, req *vo.OrderPageReq) (*comm return DefaultBillingService.GetOrderPage(userID, req) } +func (s *AccountService) ListTrafficPackages() ([]vo.TrafficPackageVO, *common.BusiError) { + return DefaultPlatformBillingService.ListActivePackages() +} + // CreateTrafficOrder 创建流量充值订单 func (s *AccountService) CreateTrafficOrder(userID int64, req *vo.TrafficOrderCreateReq) (*model.TrafficOrder, *common.BusiError) { return DefaultBillingService.CreateOrderWithPackage(userID, req.PackageID, int(req.AmountGb)) } // PayTrafficOrder 订单支付入账(admin) -func (s *AccountService) PayTrafficOrder(orderID, operatorID int64, paymentKey string) (*model.TrafficOrder, *common.BusiError) { - return DefaultBillingService.PayOrder(orderID, operatorID, paymentKey) +func (s *AccountService) CancelTrafficOrder(userID, orderID int64) *common.BusiError { + return DefaultBillingService.CancelOrder(userID, orderID) +} + +func (s *AccountService) CancelSimRechargeOrder(userID, orderID int64) *common.BusiError { + return DefaultSimRechargeService.CancelOrder(userID, orderID) } -// ListSimCards SIM 卡列表 func (s *AccountService) ListSimCards(userID int64, req *vo.SimCardPageReq) (*common.PageResponse[model.SimCard], *common.BusiError) { return DefaultBillingService.ListSimCards(userID, req) } @@ -100,6 +107,6 @@ func (s *AccountService) GetSimRechargeLogPage(userID int64, req *vo.SimRecharge } // RechargeSimCard SIM 卡充值 -func (s *AccountService) RechargeSimCard(userID, simCardID int64, req *vo.SimRechargeReq) (*model.SimRechargeLog, *common.BusiError) { - return DefaultBillingService.RechargeSimCard(userID, simCardID, req.AmountGb) +func (s *AccountService) RechargeSimCard(userID, simCardID int64, req *vo.SimRechargeReq) (*model.SimRechargeOrder, *common.BusiError) { + return DefaultSimRechargeService.CreateOrder(userID, simCardID, req.PackageID) } diff --git a/service/billing_service.go b/service/billing_service.go index f0aa4e8..edb3159 100644 --- a/service/billing_service.go +++ b/service/billing_service.go @@ -3,6 +3,7 @@ package service import ( "errors" "fmt" + "math" "strconv" "strings" "time" @@ -251,103 +252,131 @@ func (b *BillingService) CreateOrder(userID int64, amountGb int) (*model.Traffic } func (b *BillingService) CreateOrderWithPackage(userID, packageID int64, amountGb int) (*model.TrafficOrder, *common.BusiError) { - order := &model.TrafficOrder{ID: mustID(), UserID: userID, PayStatus: "unpaid", CreditStatus: "pending", CreatedAt: time.Now()} - if packageID > 0 { - var pkg model.TrafficPackage - if err := common.DB.Where("id = ? AND status = ?", packageID, "active").First(&pkg).Error; err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - return nil, common.ErrNotFound + order := &model.TrafficOrder{ID: mustID(), UserID: userID, PayStatus: "unpaid", CreditStatus: "pending", CreatedAt: time.Now(), UpdatedAt: time.Now()} + var busiErr *common.BusiError + err := common.DB.Transaction(func(tx *gorm.DB) error { + var user model.User + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Select("id").First(&user, userID).Error; err != nil { + return err + } + var pending model.TrafficOrder + if err := tx.Where("user_id = ? AND pay_status IN ?", userID, []string{"unpaid", "processing"}).First(&pending).Error; err == nil { + busiErr = common.ErrPendingOrderExists + return busiErr + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + if packageID > 0 { + var pkg model.TrafficPackage + if err := tx.Where("id = ? AND status = ?", packageID, "active").First(&pkg).Error; err != nil { + return err + } + order.PackageID, order.PackageCode = pkg.ID, pkg.Code + order.AmountBytes, order.UnitPrice, order.TotalPrice = pkg.AmountBytes, pkg.Price/float64(pkg.AmountBytes)*float64(gbBytes), pkg.Price + order.TotalFeeFen = int64(math.Round(pkg.Price * 100)) + order.AmountGb = int(pkg.AmountBytes / gbBytes) + } else { + if amountGb <= 0 { + busiErr = common.ErrParam + return busiErr } - return nil, common.ErrInternal + order.AmountGb, order.AmountBytes = amountGb, int64(amountGb)*gbBytes + order.UnitPrice, order.TotalPrice = defaultTrafficUnitPrice, float64(amountGb)*defaultTrafficUnitPrice + order.TotalFeeFen = int64(math.Round(float64(amountGb) * defaultTrafficUnitPrice * 100)) } - order.PackageID, order.PackageCode = pkg.ID, pkg.Code - order.AmountBytes, order.UnitPrice, order.TotalPrice = pkg.AmountBytes, pkg.Price/float64(pkg.AmountBytes)*float64(gbBytes), pkg.Price - order.AmountGb = int(pkg.AmountBytes / gbBytes) - } else { - if amountGb <= 0 { - return nil, common.ErrParam + return tx.Create(order).Error + }) + if err != nil { + if busiErr != nil { + return nil, busiErr + } + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, common.ErrUserNotFound + } + if code, _ := common.ParseError(err); code == 1062 { + return nil, common.ErrPendingOrderExists } - order.AmountGb, order.AmountBytes = amountGb, int64(amountGb)*gbBytes - order.UnitPrice, order.TotalPrice = defaultTrafficUnitPrice, float64(amountGb)*defaultTrafficUnitPrice - } - if err := common.DB.Create(order).Error; err != nil { logger.ERROR("创建流量订单失败", err) return nil, common.ErrInternal } + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "创建流量充值订单", fmt.Sprintf("orderId=%d amountGb=%d feeFen=%d", order.ID, order.AmountGb, order.TotalFeeFen), "success", "user_api") return order, nil } -// PayOrder 订单支付入账(admin 确认支付)。 -func (b *BillingService) PayOrder(orderID, operatorID int64, paymentKey string) (*model.TrafficOrder, *common.BusiError) { - var order model.TrafficOrder +func (b *BillingService) CancelOrder(userID, orderID int64) *common.BusiError { now := time.Now() var busiErr *common.BusiError - if paymentKey == "" { - paymentKey = strconv.FormatInt(orderID, 10) - } - paymentIdempotencyKey := "order:payment:" + paymentKey err := common.DB.Transaction(func(tx *gorm.DB) error { - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&order, orderID).Error; err != nil { + var order model.TrafficOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", orderID, userID).First(&order).Error; err != nil { return err } - if order.PayStatus == "paid" || order.CreditStatus == "credited" { - busiErr = common.ErrOrderPaid + if order.PayStatus != "unpaid" && order.PayStatus != "processing" { + busiErr = common.ErrOrderNotCancellable return busiErr } - bytes := order.AmountBytes - if bytes <= 0 { - bytes = int64(order.AmountGb) * gbBytes + result := tx.Model(&order).Where("id = ? AND pay_status IN ?", orderID, []string{"unpaid", "processing"}).Updates(map[string]any{"pay_status": "cancelled", "closed_at": now, "updated_at": now}) + if result.Error != nil { + return result.Error } - var user model.User - if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Select("id", "traffic_balance").First(&user, order.UserID).Error; err != nil { - return err + if result.RowsAffected != 1 { + busiErr = common.ErrOrderNotCancellable + return busiErr } - after := user.TrafficBalance + bytes - if err := tx.Model(&model.User{}).Where("id = ?", user.ID).UpdateColumn("traffic_balance", after).Error; err != nil { - return err + return tx.Model(&model.PaymentTransaction{}).Where("business_type = ? AND business_order_id = ? AND status IN ?", "traffic_order", orderID, []string{"unpaid", "processing"}).Updates(map[string]any{"status": "cancelled", "updated_at": now}).Error + }) + if err != nil { + if busiErr != nil { + return busiErr } - ledgerID := mustID() - if ledgerID == 0 { - return errors.New("generate ledger id failed") + if errors.Is(err, gorm.ErrRecordNotFound) { + return common.ErrOrderNotFound } - if err := tx.Create(&model.TrafficLedger{ - ID: ledgerID, AccountType: "user", AccountID: order.UserID, Direction: "credit", - AmountBytes: bytes, BalanceBefore: user.TrafficBalance, BalanceAfter: after, - SourceType: "traffic_order", SourceID: strconv.FormatInt(order.ID, 10), - IdempotencyKey: paymentIdempotencyKey, OperatorID: operatorID, CreatedAt: now, - }).Error; err != nil { - return err + logger.ERROR("取消流量订单失败", err) + return common.ErrInternal + } + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "取消流量充值订单", fmt.Sprintf("orderId=%d status=cancelled", orderID), "success", "user_api") + return nil +} + +func (b *BillingService) CloseExpiredOrders() { + cutoff := time.Now().Add(-30 * time.Minute) + var orders []model.TrafficOrder + if err := common.DB.Where("pay_status IN ? AND created_at <= ?", []string{"unpaid", "processing"}, cutoff).Limit(100).Find(&orders).Error; err != nil { + logger.ERROR("扫描超时流量订单失败", err) + } else { + for i := range orders { + b.closeTrafficOrder(orders[i].ID, cutoff) } - if err := tx.Model(&order).Updates(map[string]any{ - "pay_status": "paid", "credit_status": "credited", "amount_bytes": bytes, - "credit_ledger_id": ledgerID, "paid_at": now, "paid_by": operatorID, - "payment_idempotency_key": paymentIdempotencyKey}).Error; err != nil { + } + DefaultSimRechargeService.CloseExpiredOrders(cutoff) +} + +func (b *BillingService) closeTrafficOrder(orderID int64, cutoff time.Time) { + now := time.Now() + closed := false + var order model.TrafficOrder + err := common.DB.Transaction(func(tx *gorm.DB) error { + if err := tx.Where("id = ?", orderID).First(&order).Error; err != nil { return err } - order.AmountBytes = bytes - order.CreditStatus = "credited" - order.CreditLedgerID = ledgerID - order.PayStatus = "paid" - order.PaidAt = &now - return nil - }) - if err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - return nil, common.ErrOrderNotFound - } - if busiErr != nil { - return nil, busiErr + result := tx.Model(&model.TrafficOrder{}).Where("id = ? AND pay_status IN ? AND created_at <= ?", orderID, []string{"unpaid", "processing"}, cutoff).Updates(map[string]any{"pay_status": "closed", "closed_at": now, "updated_at": now}) + if result.Error != nil { + return result.Error } - if typed, ok := err.(*common.BusiError); ok { - return nil, typed + if result.RowsAffected != 1 { + return nil } - logger.ERROR("订单支付失败", err) - return nil, common.ErrInternal + closed = true + return tx.Model(&model.PaymentTransaction{}).Where("business_type = ? AND business_order_id = ? AND status IN ?", "traffic_order", orderID, []string{"unpaid", "processing"}).Updates(map[string]any{"status": "closed", "updated_at": now}).Error + }) + if err != nil { + logger.ERROR("关闭超时流量订单失败", err) + return } - if err := common.Delete(cache.TrafficKeyOf(order.UserID)); err != nil { - logger.WARN("删除充值余额缓存失败", err) + if closed && order.UserID != 0 { + DefaultOperationLogService.RecordEvent(order.UserID, "充值与履约", "自动关闭流量充值订单", fmt.Sprintf("orderId=%d status=closed", orderID), "success", "system") } - return &order, nil } // ListSimCards SIM 卡列表 diff --git a/service/operation_log_service.go b/service/operation_log_service.go index 8cba48e..2799568 100644 --- a/service/operation_log_service.go +++ b/service/operation_log_service.go @@ -40,6 +40,21 @@ func (s *OperationLogService) Record(userID int64, userName, module, action, det } } +func (s *OperationLogService) RecordEvent(userID int64, module, action, detail, result, source string) { + if len(detail) > 220 { + detail = detail[:220] + } + userName := "" + var user model.User + if err := common.DB.Select("name", "phone").First(&user, userID).Error; err == nil { + userName = user.Name + if userName == "" { + userName = user.Phone + } + } + s.Record(userID, userName, module, action, "source="+source+" "+detail, result, "") +} + // GetPage 操作日志分页(admin 全量,user 仅本人) func (s *OperationLogService) GetPage(req *vo.OperationLogPageReq) (*common.PageResponse[model.OperationLog], *common.BusiError) { db := common.DB.Model(&model.OperationLog{}) diff --git a/service/payment_service.go b/service/payment_service.go new file mode 100644 index 0000000..67751e3 --- /dev/null +++ b/service/payment_service.go @@ -0,0 +1,430 @@ +package service + +import ( + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "strconv" + "time" + + "gorm.io/gorm" + "gorm.io/gorm/clause" + + "laic-backend/cache" + "laic-backend/client" + "laic-backend/common" + "laic-backend/logger" + "laic-backend/model" + "laic-backend/tool" + "laic-backend/vo" +) + +type PaymentService struct { + agpay *client.AGPayClient + notifyURL string +} + +type trafficCreditAudit struct { + UserID int64 + OrderID int64 + LedgerID int64 + AmountBytes int64 + BalanceBefore int64 + BalanceAfter int64 +} + +var DefaultPaymentService = &PaymentService{} + +func InitAGPay(conf common.AGPay) error { + if conf.BaseURL == "" && conf.NotifyURL == "" { + return nil + } + instance, err := client.NewAGPayClient(conf.BaseURL) + if err != nil { + return err + } + if conf.NotifyURL == "" { + return errors.New("AGPay notify URL is required") + } + DefaultPaymentService.agpay = instance + DefaultPaymentService.notifyURL = conf.NotifyURL + return nil +} + +func (s *PaymentService) CreateWechatPayment(userID int64, req *vo.PaymentCreateReq) (*vo.PaymentCreateVO, *common.BusiError) { + if s.agpay == nil { + return nil, common.ErrPaymentDisabled + } + + transaction, created, busiErr := s.getOrCreateTransaction(userID, req.BusinessType, req.BusinessOrderID) + if busiErr != nil { + return nil, busiErr + } + if transaction.Status == "paid" { + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "创建支付交易", fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d status=paid", transaction.ID, req.BusinessType, req.BusinessOrderID), "failed", "user_api") + return nil, common.ErrPaymentAlreadyPaid + } + if !created { + if transaction.Status == "cancelled" || transaction.Status == "closed" || transaction.Status == "failed" || transaction.ProviderPayload == "" { + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "复用支付交易", fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d status=%s", transaction.ID, transaction.BusinessType, transaction.BusinessOrderID, transaction.Status), "failed", "user_api") + return nil, common.ErrPaymentOrderConflict + } + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "复用支付交易", fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d status=%s", transaction.ID, transaction.BusinessType, transaction.BusinessOrderID, transaction.Status), "success", "user_api") + return paymentCreateVO(transaction, transaction.ProviderPayload), nil + } + + payload, err := s.agpay.CreateWechatQRCode(client.AGPayCreateRequest{ + TotalFeeFen: transaction.TotalFeeFen, MerchantOrder: transaction.MerchantOrderNo, NotifyURL: s.notifyURL, + }) + if err != nil { + logger.WARN("创建 AGPay 微信支付请求失败,交易保留待确认", transaction.ID, err) + return nil, common.ErrPaymentOrderConflict + } + payloadText := string(payload) + if err := common.DB.Model(&model.PaymentTransaction{}).Where("id = ? AND status = ?", transaction.ID, "processing").Updates(map[string]any{ + "provider_payload": payloadText, "requested_at": time.Now(), "updated_at": time.Now(), + }).Error; err != nil { + logger.ERROR("保存 AGPay 支付响应失败", err) + return nil, common.ErrInternal + } + transaction.ProviderPayload = payloadText + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "创建支付交易", fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d feeFen=%d", transaction.ID, transaction.BusinessType, transaction.BusinessOrderID, transaction.TotalFeeFen), "success", "user_api") + + return paymentCreateVO(transaction, payloadText), nil +} + +func (s *PaymentService) GetTransaction(userID, transactionID int64) (*vo.PaymentTransactionVO, *common.BusiError) { + var transaction model.PaymentTransaction + if err := common.DB.Where("id = ? AND user_id = ?", transactionID, userID).First(&transaction).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, common.ErrPaymentTransactionNotFound + } + return nil, common.ErrInternal + } + return paymentTransactionVO(&transaction), nil +} + +func (s *PaymentService) markPaymentRequestFailed(transactionID int64, businessType string, businessOrderID int64) { + if err := common.DB.Transaction(func(tx *gorm.DB) error { + if err := tx.Model(&model.PaymentTransaction{}).Where("id = ? AND status = ?", transactionID, "processing").Updates(map[string]any{"status": "failed", "updated_at": time.Now()}).Error; err != nil { + return err + } + switch businessType { + case "traffic_order": + return tx.Model(&model.TrafficOrder{}).Where("id = ? AND pay_status = ?", businessOrderID, "processing").Updates(map[string]any{"pay_status": "unpaid", "payment_provider": "", "updated_at": time.Now()}).Error + case "sim_recharge_order": + return tx.Model(&model.SimRechargeOrder{}).Where("id = ? AND payment_status = ?", businessOrderID, "processing").Updates(map[string]any{"payment_status": "unpaid", "updated_at": time.Now()}).Error + default: + return common.ErrParam + } + }); err != nil { + logger.ERROR("回滚支付请求状态失败", err) + } +} + +func (s *PaymentService) getOrCreateTransaction(userID int64, businessType string, businessOrderID int64) (*model.PaymentTransaction, bool, *common.BusiError) { + var transaction model.PaymentTransaction + var created bool + var busiErr *common.BusiError + err := common.DB.Transaction(func(tx *gorm.DB) error { + amount, status, err := paymentBusinessOrder(tx, userID, businessType, businessOrderID) + if err != nil { + return err + } + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("user_id = ? AND business_type = ? AND business_order_id = ? AND provider = ?", userID, businessType, businessOrderID, "agpay"). + First(&transaction).Error; err == nil { + return nil + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + if amount <= 0 { + busiErr = common.ErrParam + return busiErr + } + if status == "paid" { + busiErr = common.ErrPaymentAlreadyPaid + return busiErr + } + if status != "unpaid" && status != "processing" { + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + + id, err := tool.NextID() + if err != nil { + return err + } + now := time.Now() + transaction = model.PaymentTransaction{ + ID: id, UserID: userID, BusinessType: businessType, BusinessOrderID: businessOrderID, + Provider: "agpay", MerchantOrderNo: "agpay-" + strconv.FormatInt(id, 10), + TotalFeeFen: amount, Status: "processing", CreatedAt: now, UpdatedAt: now, + } + if err := tx.Create(&transaction).Error; err != nil { + return err + } + created = true + return updatePaymentBusinessStatus(tx, businessType, businessOrderID, "processing") + }) + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, false, common.ErrOrderNotFound + } + if busiErr != nil { + return nil, false, busiErr + } + logger.ERROR("创建支付交易失败", err) + return nil, false, common.ErrInternal + } + return &transaction, created, nil +} + +func paymentBusinessOrder(tx *gorm.DB, userID int64, businessType string, businessOrderID int64) (int64, string, error) { + switch businessType { + case "traffic_order": + var order model.TrafficOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", businessOrderID, userID).First(&order).Error; err != nil { + return 0, "", err + } + return order.TotalFeeFen, order.PayStatus, nil + case "sim_recharge_order": + var order model.SimRechargeOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", businessOrderID, userID).First(&order).Error; err != nil { + return 0, "", err + } + return order.TotalFeeFen, order.PaymentStatus, nil + default: + return 0, "", common.ErrParam + } +} + +func updatePaymentBusinessStatus(tx *gorm.DB, businessType string, businessOrderID int64, status string) error { + var result *gorm.DB + switch businessType { + case "traffic_order": + result = tx.Model(&model.TrafficOrder{}).Where("id = ? AND pay_status = ?", businessOrderID, "unpaid").Updates(map[string]any{"pay_status": status, "payment_provider": "agpay", "updated_at": time.Now()}) + case "sim_recharge_order": + result = tx.Model(&model.SimRechargeOrder{}).Where("id = ? AND payment_status = ?", businessOrderID, "unpaid").Updates(map[string]any{"payment_status": status, "updated_at": time.Now()}) + default: + return common.ErrParam + } + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return common.ErrPaymentOrderConflict + } + return nil +} + +func paymentCreateVO(transaction *model.PaymentTransaction, payload string) *vo.PaymentCreateVO { + return &vo.PaymentCreateVO{ + TransactionID: transaction.ID, BusinessType: transaction.BusinessType, BusinessOrderID: transaction.BusinessOrderID, + Provider: transaction.Provider, MerchantOrderNo: transaction.MerchantOrderNo, TotalFeeFen: transaction.TotalFeeFen, + Status: transaction.Status, ProviderPayload: payload, CallbackEnabled: false, + } +} + +func paymentTransactionVO(transaction *model.PaymentTransaction) *vo.PaymentTransactionVO { + return &vo.PaymentTransactionVO{ + ID: transaction.ID, BusinessType: transaction.BusinessType, BusinessOrderID: transaction.BusinessOrderID, + Provider: transaction.Provider, MerchantOrderNo: transaction.MerchantOrderNo, TotalFeeFen: transaction.TotalFeeFen, + Status: transaction.Status, CreatedAt: transaction.CreatedAt, PaidAt: transaction.PaidAt, + } +} + +func (s *PaymentService) ConfirmAGPayCallback(tradeNo string, payload []byte) *common.BusiError { + if tradeNo == "" || len(tradeNo) > 128 { + return common.ErrParam + } + + now := time.Now() + sum := sha256.Sum256(payload) + payloadHash := hex.EncodeToString(sum[:]) + creditedUserID := int64(0) + var trafficCredit *trafficCreditAudit + var simOrderToSubmit *model.SimRechargeOrder + callbackAudit := "" + var callbackTransaction model.PaymentTransaction + var busiErr *common.BusiError + err := common.DB.Transaction(func(tx *gorm.DB) error { + var transaction model.PaymentTransaction + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("provider = ? AND merchant_order_no = ?", "agpay", tradeNo).First(&transaction).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + busiErr = common.ErrPaymentTransactionNotFound + return busiErr + } + return err + } + + var event model.PaymentCallbackEvent + err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("provider = ? AND event_key = ?", "agpay", tradeNo).First(&event).Error + if err == nil { + callbackTransaction = transaction + if event.PayloadHash != payloadHash { + callbackAudit = "conflict" + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + callbackAudit = "duplicate" + callbackTransaction = transaction + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + id, err := tool.NextID() + if err != nil { + return err + } + event = model.PaymentCallbackEvent{ID: id, Provider: "agpay", EventKey: tradeNo, PayloadHash: payloadHash, Verified: true, CreatedAt: now} + if err := tx.Create(&event).Error; err != nil { + return err + } + + callbackTransaction = transaction + if transaction.Status == "paid" { + callbackAudit = "duplicate" + return tx.Model(&event).Updates(map[string]any{"processed_at": now, "result_code": "duplicate"}).Error + } + if transaction.Status == "cancelled" || transaction.Status == "closed" { + callbackAudit = transaction.Status + return tx.Model(&event).Updates(map[string]any{"processed_at": now, "result_code": transaction.Status}).Error + } + if transaction.Status != "processing" { + callbackAudit = "conflict" + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + result := tx.Model(&transaction).Where("id = ? AND status = ?", transaction.ID, "processing").Updates(map[string]any{ + "status": "paid", "paid_at": now, "updated_at": now, + }) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + callbackAudit = "conflict" + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + switch transaction.BusinessType { + case "traffic_order": + trafficCredit, err = creditPaidTrafficOrder(tx, transaction.BusinessOrderID, transaction.UserID, now) + if err != nil { + return err + } + creditedUserID = trafficCredit.UserID + case "sim_recharge_order": + var order model.SimRechargeOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", transaction.BusinessOrderID, transaction.UserID).First(&order).Error; err != nil { + return err + } + result := tx.Model(&model.SimRechargeOrder{}).Where("id = ? AND user_id = ? AND payment_status = ? AND fulfillment_status = ?", order.ID, transaction.UserID, "processing", "pending").Updates(map[string]any{ + "payment_status": "paid", "paid_at": now, "fulfillment_status": "submitting", "attempt_count": 1, + "fulfillment_started_at": now, "fulfillment_message": "续费提交处理中;如状态未更新,请联系管理员。", + "fulfillment_lease_token": nil, "fulfillment_lease_until": nil, "next_attempt_at": nil, + "last_error_code": "", "last_error_message": "", "updated_at": now, + }) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + callbackAudit = "conflict" + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + order.PaymentStatus = "paid" + order.FulfillmentStatus = "submitting" + order.AttemptCount = 1 + order.PaidAt = &now + order.FulfillmentStartedAt = &now + order.FulfillmentMessage = "续费提交处理中;如状态未更新,请联系管理员。" + simOrderToSubmit = &order + default: + callbackAudit = "conflict" + busiErr = common.ErrPaymentOrderConflict + return busiErr + } + callbackAudit = "paid" + return tx.Model(&event).Updates(map[string]any{"processed_at": now, "result_code": "paid"}).Error + }) + if err != nil { + if callbackAudit == "conflict" && callbackTransaction.ID != 0 { + DefaultOperationLogService.RecordEvent(callbackTransaction.UserID, "充值与履约", "AGPay回调冲突", fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d", callbackTransaction.ID, callbackTransaction.BusinessType, callbackTransaction.BusinessOrderID), "failed", "agpay_callback") + } + if busiErr != nil { + return busiErr + } + logger.ERROR("确认 AGPay 支付回调失败", err) + return common.ErrInternal + } + if callbackAudit != "" { + action := "确认AGPay支付回调" + if callbackAudit == "duplicate" { + action = "忽略重复AGPay回调" + } else if callbackAudit == "cancelled" || callbackAudit == "closed" { + action = "忽略迟到AGPay回调" + } + DefaultOperationLogService.RecordEvent(callbackTransaction.UserID, "充值与履约", action, fmt.Sprintf("transactionId=%d businessType=%s businessOrderId=%d result=%s", callbackTransaction.ID, callbackTransaction.BusinessType, callbackTransaction.BusinessOrderID, callbackAudit), "success", "agpay_callback") + } + if trafficCredit != nil { + DefaultOperationLogService.RecordEvent(trafficCredit.UserID, "充值与履约", "流量充值入账", fmt.Sprintf("orderId=%d ledgerId=%d amountBytes=%d balanceBefore=%d balanceAfter=%d", trafficCredit.OrderID, trafficCredit.LedgerID, trafficCredit.AmountBytes, trafficCredit.BalanceBefore, trafficCredit.BalanceAfter), "success", "agpay_callback") + } + if creditedUserID != 0 { + if err := common.Delete(cache.TrafficKeyOf(creditedUserID)); err != nil { + logger.WARN("删除支付入账余额缓存失败", err) + } + } + if simOrderToSubmit != nil { + DefaultOperationLogService.RecordEvent(simOrderToSubmit.UserID, "充值与履约", "启动SIM续费履约", fmt.Sprintf("orderId=%d attempt=1", simOrderToSubmit.ID), "success", "agpay_callback") + DefaultSimRechargeService.SubmitPaidOrderOnce(simOrderToSubmit) + } + return nil +} + +func creditPaidTrafficOrder(tx *gorm.DB, orderID, expectedUserID int64, now time.Time) (*trafficCreditAudit, error) { + var order model.TrafficOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", orderID, expectedUserID).First(&order).Error; err != nil { + return nil, err + } + if order.PayStatus == "paid" && order.CreditStatus == "credited" { + return nil, nil + } + if order.PayStatus != "processing" || order.CreditStatus != "pending" || order.AmountBytes <= 0 { + return nil, common.ErrPaymentOrderConflict + } + var user model.User + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Select("id", "traffic_balance").First(&user, order.UserID).Error; err != nil { + return nil, err + } + ledgerID := mustID() + if ledgerID == 0 { + return nil, errors.New("generate traffic credit ledger ID") + } + after := user.TrafficBalance + order.AmountBytes + if err := tx.Model(&user).UpdateColumn("traffic_balance", after).Error; err != nil { + return nil, err + } + idempotencyKey := "payment:" + strconv.FormatInt(order.ID, 10) + if err := tx.Create(&model.TrafficLedger{ + ID: ledgerID, AccountType: "user", AccountID: order.UserID, Direction: "credit", AmountBytes: order.AmountBytes, + BalanceBefore: user.TrafficBalance, BalanceAfter: after, SourceType: "traffic_order", SourceID: strconv.FormatInt(order.ID, 10), + IdempotencyKey: idempotencyKey, CreatedAt: now, + }).Error; err != nil { + return nil, err + } + result := tx.Model(&model.TrafficOrder{}).Where("id = ? AND user_id = ? AND pay_status = ? AND credit_status = ?", order.ID, expectedUserID, "processing", "pending").Updates(map[string]any{ + "pay_status": "paid", "payment_provider": "agpay", + "credit_status": "credited", "credit_ledger_id": ledgerID, "payment_idempotency_key": idempotencyKey, + "paid_at": now, "updated_at": now, + }) + if result.Error != nil { + return nil, result.Error + } + if result.RowsAffected != 1 { + return nil, common.ErrPaymentOrderConflict + } + return &trafficCreditAudit{UserID: order.UserID, OrderID: order.ID, LedgerID: ledgerID, AmountBytes: order.AmountBytes, BalanceBefore: user.TrafficBalance, BalanceAfter: after}, nil +} diff --git a/service/platform_billing_service.go b/service/platform_billing_service.go index 9cea616..72b90a8 100644 --- a/service/platform_billing_service.go +++ b/service/platform_billing_service.go @@ -154,6 +154,18 @@ func (s *PlatformBillingService) GetPackagePage(req *vo.TrafficPackagePageReq) ( return common.Page(req.Pagination, total, list), nil } +func (s *PlatformBillingService) ListActivePackages() ([]vo.TrafficPackageVO, *common.BusiError) { + var packages []vo.TrafficPackageVO + if err := common.DB.Model(&model.TrafficPackage{}). + Select("id, code, name, amount_bytes, price, validity_days, sort"). + Where("status = ?", "active"). + Order("sort ASC, id DESC"). + Find(&packages).Error; err != nil { + return nil, common.ErrInternal + } + return packages, nil +} + func (s *PlatformBillingService) UpdatePackage(id int64, req *vo.TrafficPackageUpdateReq) *common.BusiError { updates := map[string]any{"updated_at": time.Now(), "sort": req.Sort} if req.Name != "" { diff --git a/service/scheduler.go b/service/scheduler.go index 2beb5c5..30b14f2 100644 --- a/service/scheduler.go +++ b/service/scheduler.go @@ -42,6 +42,9 @@ func (s *Scheduler) Start() { _, _ = s.cron.AddFunc("@every 30m", DefaultBillingService.flushBalanceSnapshot) // SIM 卡用量:每 1h 从运营商同步用量 _, _ = s.cron.AddFunc("@every 1h", DefaultBillingService.syncSimUsage) + _, _ = s.cron.AddFunc("@every 30s", func() { + DefaultBillingService.CloseExpiredOrders() + }) // 设备心跳:清理未再收到真实 MQTT 消息的在线设备 _, _ = s.cron.AddFunc("@every "+heartbeatScanInterval().String(), expireDeviceHeartbeats) diff --git a/service/sim_recharge_service.go b/service/sim_recharge_service.go new file mode 100644 index 0000000..f147023 --- /dev/null +++ b/service/sim_recharge_service.go @@ -0,0 +1,324 @@ +package service + +import ( + "errors" + "fmt" + "strconv" + "time" + + "gorm.io/gorm" + "gorm.io/gorm/clause" + + "laic-backend/client" + "laic-backend/common" + "laic-backend/logger" + "laic-backend/model" + "laic-backend/tool" +) + +type simbossAdapter struct { + client *client.SimbossClient +} + +func InitSimboss(conf common.Simboss) error { + if conf.AppID == "" && conf.Secret == "" { + return nil + } + instance, err := client.NewSimbossClient(conf.AppID, conf.Secret, conf.APIBase) + if err != nil { + return err + } + DefaultSimRechargeService.simboss = &simbossAdapter{client: instance} + return nil +} + +type SimRechargeService struct { + simboss *simbossAdapter +} + +var DefaultSimRechargeService = &SimRechargeService{} + +func (s *SimRechargeService) GetCard(userID, cardID int64, refresh bool) (*model.SimCard, *common.BusiError) { + var card model.SimCard + if err := common.DB.Where("id = ? AND user_id = ?", cardID, userID).First(&card).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, common.ErrSimCardNotFound + } + return nil, common.ErrInternal + } + if refresh { + if s.simboss == nil { + return nil, common.ErrCarrierUnavailable + } + if err := s.refreshCard(&card); err != nil { + logger.WARN("刷新 SIM 卡信息失败", card.ID, err) + return nil, common.ErrCarrierUnavailable + } + } + return &card, nil +} + +func (s *SimRechargeService) ListPackages() ([]model.SimPackage, *common.BusiError) { + var packages []model.SimPackage + if err := common.DB.Where("status = ?", "active").Order("sale_fee_fen ASC, id ASC").Find(&packages).Error; err != nil { + return nil, common.ErrInternal + } + return packages, nil +} + +func (s *SimRechargeService) CreateOrder(userID, cardID, packageID int64) (*model.SimRechargeOrder, *common.BusiError) { + var order model.SimRechargeOrder + now := time.Now() + var busiErr *common.BusiError + err := common.DB.Transaction(func(tx *gorm.DB) error { + var card model.SimCard + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", cardID, userID).First(&card).Error; err != nil { + return err + } + var pending model.SimRechargeOrder + if err := tx.Where("sim_card_id = ? AND payment_status IN ?", cardID, []string{"unpaid", "processing"}).First(&pending).Error; err == nil { + busiErr = common.ErrPendingOrderExists + return busiErr + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + if card.Status != "active" || card.CarrierStatus == "cancelled" || card.CarrierStatus == "suspended" || card.Iccid == "" { + busiErr = common.ErrParam + return busiErr + } + var pkg model.SimPackage + if err := tx.Where("id = ? AND status = ?", packageID, "active").First(&pkg).Error; err != nil { + return err + } + if pkg.ProviderRatePlanID <= 0 || pkg.AmountGb <= 0 || pkg.ValidityMonths <= 0 || pkg.SaleFeeFen <= 0 { + busiErr = common.ErrParam + return busiErr + } + if card.Carrier != "" && pkg.Carrier != "" && card.Carrier != pkg.Carrier { + busiErr = common.ErrParam + return busiErr + } + id, err := tool.NextID() + if err != nil { + return err + } + order = model.SimRechargeOrder{ID: id, UserID: userID, SimCardID: card.ID, PackageID: pkg.ID, PackageCode: pkg.Code, ProviderRatePlanID: pkg.ProviderRatePlanID, IccidSnapshot: card.Iccid, AmountGb: pkg.AmountGb, ValidityMonths: pkg.ValidityMonths, TotalFeeFen: pkg.SaleFeeFen, PaymentStatus: "unpaid", FulfillmentStatus: "pending", ExternalOrderNo: "sim-" + strconv.FormatInt(id, 10), CreatedAt: now, UpdatedAt: now} + return tx.Create(&order).Error + }) + if err != nil { + if busiErr != nil { + return nil, busiErr + } + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, common.ErrSimCardNotFound + } + if code, _ := common.ParseError(err); code == 1062 { + return nil, common.ErrPendingOrderExists + } + logger.ERROR("创建 SIM 续费订单失败", err) + return nil, common.ErrInternal + } + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "创建SIM续费订单", fmt.Sprintf("orderId=%d simCardId=%d feeFen=%d", order.ID, order.SimCardID, order.TotalFeeFen), "success", "user_api") + return &order, nil +} + +func (s *SimRechargeService) CloseExpiredOrders(cutoff time.Time) { + var orders []model.SimRechargeOrder + if err := common.DB.Where("payment_status IN ? AND created_at <= ?", []string{"unpaid", "processing"}, cutoff).Limit(100).Find(&orders).Error; err != nil { + logger.ERROR("扫描超时 SIM 订单失败", err) + return + } + for i := range orders { + now := time.Now() + closed := false + err := common.DB.Transaction(func(tx *gorm.DB) error { + result := tx.Model(&model.SimRechargeOrder{}).Where("id = ? AND payment_status IN ? AND created_at <= ?", orders[i].ID, []string{"unpaid", "processing"}, cutoff).Updates(map[string]any{"payment_status": "closed", "closed_at": now, "updated_at": now}) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return nil + } + closed = true + return tx.Model(&model.PaymentTransaction{}).Where("business_type = ? AND business_order_id = ? AND status IN ?", "sim_recharge_order", orders[i].ID, []string{"unpaid", "processing"}).Updates(map[string]any{"status": "closed", "updated_at": now}).Error + }) + if err != nil { + logger.ERROR("关闭超时 SIM 订单失败", err) + continue + } + if closed { + DefaultOperationLogService.RecordEvent(orders[i].UserID, "充值与履约", "自动关闭SIM续费订单", fmt.Sprintf("orderId=%d status=closed", orders[i].ID), "success", "system") + } + } +} + +func (s *SimRechargeService) CancelOrder(userID, orderID int64) *common.BusiError { + now := time.Now() + var busiErr *common.BusiError + err := common.DB.Transaction(func(tx *gorm.DB) error { + var order model.SimRechargeOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ? AND user_id = ?", orderID, userID).First(&order).Error; err != nil { + return err + } + if order.PaymentStatus != "unpaid" && order.PaymentStatus != "processing" { + busiErr = common.ErrOrderNotCancellable + return busiErr + } + result := tx.Model(&order).Where("id = ? AND payment_status IN ?", orderID, []string{"unpaid", "processing"}).Updates(map[string]any{"payment_status": "cancelled", "closed_at": now, "updated_at": now}) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + busiErr = common.ErrOrderNotCancellable + return busiErr + } + return tx.Model(&model.PaymentTransaction{}).Where("business_type = ? AND business_order_id = ? AND status IN ?", "sim_recharge_order", orderID, []string{"unpaid", "processing"}).Updates(map[string]any{"status": "cancelled", "updated_at": now}).Error + }) + if err != nil { + if busiErr != nil { + return busiErr + } + if errors.Is(err, gorm.ErrRecordNotFound) { + return common.ErrOrderNotFound + } + logger.ERROR("取消 SIM 订单失败", err) + return common.ErrInternal + } + DefaultOperationLogService.RecordEvent(userID, "充值与履约", "取消SIM续费订单", fmt.Sprintf("orderId=%d status=cancelled", orderID), "success", "user_api") + return nil +} + +func (s *SimRechargeService) GetOrder(userID, orderID int64) (*model.SimRechargeOrder, *common.BusiError) { + var order model.SimRechargeOrder + if err := common.DB.Where("id = ? AND user_id = ?", orderID, userID).First(&order).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, common.ErrOrderNotFound + } + return nil, common.ErrInternal + } + return &order, nil +} + +func (s *SimRechargeService) SubmitPaidOrderOnce(order *model.SimRechargeOrder) { + const supportMessage = "续费处理异常,请联系管理员。" + DefaultOperationLogService.RecordEvent(order.UserID, "充值与履约", "提交SIMBOSS履约", fmt.Sprintf("orderId=%d attempt=1", order.ID), "success", "simboss") + if s.simboss == nil { + DefaultOperationLogService.RecordEvent(order.UserID, "充值与履约", "SIMBOSS履约转人工处理", fmt.Sprintf("orderId=%d errorCode=simboss_unavailable", order.ID), "failed", "simboss") + s.markManualReview(order.ID, "simboss_unavailable", supportMessage, errors.New("SIMBOSS is not configured")) + return + } + providerOrderNo, err := s.simboss.client.Recharge(order.IccidSnapshot, order.ProviderRatePlanID, order.ValidityMonths, order.ExternalOrderNo) + if err != nil { + logger.ERROR("SIM 续费提交 SIMBOSS 失败", err) + DefaultOperationLogService.RecordEvent(order.UserID, "充值与履约", "SIMBOSS履约转人工处理", fmt.Sprintf("orderId=%d errorCode=simboss_request_failed", order.ID), "failed", "simboss") + s.markManualReview(order.ID, "simboss_request_failed", supportMessage, err) + return + } + if providerOrderNo == "" { + err = errors.New("SIMBOSS returned an empty order number") + logger.ERROR("SIMBOSS 返回空订单号", err) + DefaultOperationLogService.RecordEvent(order.UserID, "充值与履约", "SIMBOSS履约转人工处理", fmt.Sprintf("orderId=%d errorCode=simboss_empty_order_no", order.ID), "failed", "simboss") + s.markManualReview(order.ID, "simboss_empty_order_no", supportMessage, err) + return + } + now := time.Now() + result := common.DB.Model(&model.SimRechargeOrder{}).Where("id = ? AND payment_status = ? AND fulfillment_status = ? AND attempt_count = ?", order.ID, "paid", "submitting", 1).Updates(map[string]any{ + "fulfillment_status": "success", "provider_order_no": providerOrderNo, "fulfilled_at": now, + "fulfillment_message": "", "last_error_code": "", "last_error_message": "", "updated_at": now, + }) + if result.Error != nil { + logger.ERROR("更新 SIM 续费履约结果失败", result.Error) + return + } + if result.RowsAffected != 1 { + logger.WARN("SIM 续费履约状态已变化", order.ID, errors.New("stale fulfillment state")) + } +} + +func (s *SimRechargeService) markManualReview(orderID int64, code, message string, cause error) { + now := time.Now() + result := common.DB.Model(&model.SimRechargeOrder{}).Where("id = ? AND payment_status = ? AND fulfillment_status = ? AND attempt_count = ?", orderID, "paid", "submitting", 1).Updates(map[string]any{ + "fulfillment_status": "manual_review", "fulfillment_message": message, + "last_error_code": code, "last_error_message": cause.Error(), "next_attempt_at": nil, + "fulfillment_lease_token": nil, "fulfillment_lease_until": nil, "updated_at": now, + }) + if result.Error != nil { + logger.ERROR("标记 SIM 续费人工处理失败", result.Error) + return + } + if result.RowsAffected != 1 { + logger.WARN("SIM 续费人工处理状态已变化", orderID, errors.New("stale fulfillment state")) + } +} + +func (s *SimRechargeService) refreshCard(card *model.SimCard) error { + detail, err := s.simboss.client.GetDeviceDetail(card.Iccid) + if err != nil { + return err + } + now := time.Now() + updates := map[string]any{ + "carrier": detail.Carrier, "carrier_status": normalizeCarrierStatus(detail.Status, detail.DeviceStatus), + "used_gb": detail.DataUsage / 1024, "last_sync_at": now, "updated_at": now, + } + if detail.RatePlanExpirationDate != "" { + if expiresAt, err := parseSimbossDate(detail.RatePlanExpirationDate); err == nil { + updates["expired_at"] = expiresAt + } + } + if updates["carrier_status"] == "cancelled" { + updates["status"] = "expired" + } + if err := common.DB.Model(card).Updates(updates).Error; err != nil { + return err + } + return common.DB.First(card, card.ID).Error +} + +func parseSimbossDate(value string) (time.Time, error) { + for _, layout := range []string{"2006-01-02", time.RFC3339, "2006-01-02 15:04:05"} { + if parsed, err := time.ParseInLocation(layout, value, time.Local); err == nil { + return parsed, nil + } + } + return time.Time{}, errors.New("invalid SIMBOSS date") +} + +func normalizeCarrierStatus(status, deviceStatus string) string { + value := status + " " + deviceStatus + switch { + case containsFold(value, "cancel"): + return "cancelled" + case containsFold(value, "arrear"): + return "arrears" + case containsFold(value, "suspend"): + return "suspended" + default: + return "normal" + } +} + +func containsFold(value, fragment string) bool { + for i := 0; i+len(fragment) <= len(value); i++ { + if equalFoldASCII(value[i:i+len(fragment)], fragment) { + return true + } + } + return false +} + +func equalFoldASCII(value, fragment string) bool { + for i := range value { + left, right := value[i], fragment[i] + if left >= 'A' && left <= 'Z' { + left += 'a' - 'A' + } + if right >= 'A' && right <= 'Z' { + right += 'a' - 'A' + } + if left != right { + return false + } + } + return true +} diff --git a/sql/014_payment_and_sim_recharge.sql b/sql/014_payment_and_sim_recharge.sql new file mode 100644 index 0000000..786ac7e --- /dev/null +++ b/sql/014_payment_and_sim_recharge.sql @@ -0,0 +1,91 @@ +-- 支付与 SIM 套餐续费基础。执行前请先备份生产数据库。 + +ALTER TABLE traffic_order + ADD COLUMN total_fee_fen BIGINT NOT NULL DEFAULT 0 COMMENT '固定支付金额(分)' AFTER total_price, + ADD COLUMN payment_provider VARCHAR(32) DEFAULT '' COMMENT '支付渠道' AFTER pay_status, + ADD COLUMN provider_trade_no VARCHAR(128) DEFAULT '' COMMENT '第三方交易号' AFTER payment_provider, + ADD COLUMN credit_attempt_count INT NOT NULL DEFAULT 0 COMMENT '入账尝试次数' AFTER credit_status, + ADD COLUMN next_credit_at DATETIME NULL COMMENT '下次入账时间' AFTER credit_attempt_count, + ADD COLUMN closed_at DATETIME NULL COMMENT '关闭时间' AFTER paid_at, + ADD UNIQUE KEY uk_traffic_order_provider_trade (payment_provider, provider_trade_no), + ADD KEY idx_traffic_order_credit_retry (pay_status, credit_status, next_credit_at); + +UPDATE traffic_order +SET total_fee_fen = ROUND(total_price * 100) +WHERE total_fee_fen = 0 AND total_price IS NOT NULL; + +CREATE TABLE IF NOT EXISTS sim_package ( + id BIGINT PRIMARY KEY COMMENT '主键 ID', + code VARCHAR(64) NOT NULL COMMENT '平台套餐编码', + name VARCHAR(128) NOT NULL COMMENT '套餐名称', + carrier VARCHAR(16) NOT NULL COMMENT '运营商', + provider_rate_plan_id BIGINT NOT NULL COMMENT 'SIMBOSS 套餐标识', + amount_gb INT NOT NULL COMMENT '流量额度 GB', + validity_months INT NOT NULL COMMENT '有效月数', + cost_fee_fen BIGINT NOT NULL COMMENT '采购成本(分)', + sale_fee_fen BIGINT NOT NULL COMMENT '销售价格(分)', + status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive', + created_at DATETIME NULL COMMENT '创建时间', + updated_at DATETIME NULL COMMENT '更新时间', + UNIQUE KEY uk_sim_package_code (code), + UNIQUE KEY uk_sim_package_provider_plan (carrier, provider_rate_plan_id), + KEY idx_sim_package_status (status) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE IF NOT EXISTS sim_recharge_order ( + id BIGINT PRIMARY KEY COMMENT '主键 ID', + user_id BIGINT NOT NULL COMMENT '所属用户 ID', + sim_card_id BIGINT NOT NULL COMMENT 'SIM 卡 ID', + package_id BIGINT NOT NULL COMMENT '套餐 ID', + package_code VARCHAR(64) NOT NULL COMMENT '套餐编码快照', + provider_rate_plan_id BIGINT NOT NULL COMMENT 'SIMBOSS 套餐标识快照', + iccid_snapshot VARCHAR(32) NOT NULL COMMENT 'ICCID 快照,仅供履约', + amount_gb INT NOT NULL COMMENT '流量额度 GB', + validity_months INT NOT NULL COMMENT '有效月数', + total_fee_fen BIGINT NOT NULL COMMENT '固定支付金额(分)', + payment_status VARCHAR(16) NOT NULL DEFAULT 'unpaid' COMMENT 'unpaid/processing/paid/failed/closed', + fulfillment_status VARCHAR(16) NOT NULL DEFAULT 'pending' COMMENT 'pending/submitting/retrying/success/failed', + external_order_no VARCHAR(128) NOT NULL COMMENT '平台稳定外部订单号', + provider_order_no VARCHAR(128) DEFAULT '' COMMENT '运营商流水号', + attempt_count INT NOT NULL DEFAULT 0 COMMENT '履约尝试次数', + next_attempt_at DATETIME NULL COMMENT '下次履约时间', + last_error_code VARCHAR(64) DEFAULT '' COMMENT '最后错误代码', + last_error_message VARCHAR(256) DEFAULT '' COMMENT '最后错误信息', + created_at DATETIME NULL COMMENT '创建时间', + paid_at DATETIME NULL COMMENT '支付时间', + fulfilled_at DATETIME NULL COMMENT '履约时间', + updated_at DATETIME NULL COMMENT '更新时间', + UNIQUE KEY uk_sim_recharge_external (external_order_no), + KEY idx_sim_recharge_user_created (user_id, created_at), + KEY idx_sim_recharge_fulfillment_retry (payment_status, fulfillment_status, next_attempt_at) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE IF NOT EXISTS payment_transaction ( + id BIGINT PRIMARY KEY COMMENT '主键 ID', + user_id BIGINT NOT NULL COMMENT '付款用户 ID', + business_type VARCHAR(32) NOT NULL COMMENT 'traffic_order/sim_recharge_order', + business_order_id BIGINT NOT NULL COMMENT '业务订单 ID', + provider VARCHAR(32) NOT NULL COMMENT '支付渠道', + merchant_order_no VARCHAR(128) NOT NULL COMMENT '商户订单号', + provider_trade_no VARCHAR(128) DEFAULT '' COMMENT '第三方交易号', + total_fee_fen BIGINT NOT NULL COMMENT '固定金额(分)', + status VARCHAR(16) NOT NULL DEFAULT 'unpaid' COMMENT 'unpaid/processing/paid/failed/closed', + created_at DATETIME NULL COMMENT '创建时间', + paid_at DATETIME NULL COMMENT '支付时间', + updated_at DATETIME NULL COMMENT '更新时间', + UNIQUE KEY uk_payment_merchant_order (provider, merchant_order_no), + UNIQUE KEY uk_payment_provider_trade (provider, provider_trade_no), + KEY idx_payment_business (business_type, business_order_id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + +CREATE TABLE IF NOT EXISTS payment_callback_event ( + id BIGINT PRIMARY KEY COMMENT '主键 ID', + provider VARCHAR(32) NOT NULL COMMENT '支付渠道', + event_key VARCHAR(128) NOT NULL COMMENT '回调事件幂等键', + payload_hash VARCHAR(64) NOT NULL COMMENT '载荷哈希', + verified TINYINT(1) NOT NULL DEFAULT 0 COMMENT '是否通过验签', + processed_at DATETIME NULL COMMENT '处理时间', + result_code VARCHAR(64) DEFAULT '' COMMENT '处理结果', + created_at DATETIME NULL COMMENT '创建时间', + UNIQUE KEY uk_payment_callback_event (provider, event_key) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/sql/015_sim_recharge_order_rate_plan_snapshot.sql b/sql/015_sim_recharge_order_rate_plan_snapshot.sql new file mode 100644 index 0000000..66dd651 --- /dev/null +++ b/sql/015_sim_recharge_order_rate_plan_snapshot.sql @@ -0,0 +1,12 @@ +-- 为已存在的 sim_recharge_order 表补充套餐快照字段。 +-- 新部署使用 014_payment_and_sim_recharge.sql 建表时已包含此列,无需重复执行本文件。 +ALTER TABLE sim_recharge_order + ADD COLUMN provider_rate_plan_id BIGINT NOT NULL DEFAULT 0 COMMENT 'SIMBOSS 套餐标识快照' AFTER package_code; + +UPDATE sim_recharge_order o +LEFT JOIN sim_package p ON p.id = o.package_id +SET o.provider_rate_plan_id = COALESCE(p.provider_rate_plan_id, 0) +WHERE o.provider_rate_plan_id = 0; + +ALTER TABLE sim_recharge_order + ALTER COLUMN provider_rate_plan_id DROP DEFAULT; diff --git a/sql/016_payment_and_fulfillment_hardening.sql b/sql/016_payment_and_fulfillment_hardening.sql new file mode 100644 index 0000000..19fcb49 --- /dev/null +++ b/sql/016_payment_and_fulfillment_hardening.sql @@ -0,0 +1,22 @@ +-- 支付与 SIM 履约加固。执行前请备份数据库,并先检查现有重复数据。 + +UPDATE payment_transaction SET provider_trade_no = NULL WHERE provider_trade_no = ''; +UPDATE traffic_order SET provider_trade_no = NULL WHERE provider_trade_no = ''; + +ALTER TABLE payment_transaction + DROP INDEX uk_payment_provider_trade, + MODIFY COLUMN provider_trade_no VARCHAR(128) NULL DEFAULT NULL COMMENT '第三方交易号', + ADD COLUMN provider_payload TEXT NULL COMMENT 'AGPay下单响应快照' AFTER provider_trade_no, + ADD COLUMN requested_at DATETIME NULL COMMENT '最近一次请求AGPay时间' AFTER status, + ADD UNIQUE KEY uk_payment_provider_trade (provider, provider_trade_no), + ADD UNIQUE KEY uk_payment_business (provider, business_type, business_order_id); + +ALTER TABLE traffic_order + DROP INDEX uk_traffic_order_provider_trade, + MODIFY COLUMN provider_trade_no VARCHAR(128) NULL DEFAULT NULL COMMENT '第三方交易号', + ADD UNIQUE KEY uk_traffic_order_provider_trade (payment_provider, provider_trade_no); + +ALTER TABLE sim_recharge_order + ADD COLUMN fulfillment_lease_token VARCHAR(64) NULL COMMENT '履约租约令牌' AFTER fulfillment_status, + ADD COLUMN fulfillment_lease_until DATETIME NULL COMMENT '履约租约截止时间' AFTER fulfillment_lease_token, + ADD KEY idx_sim_recharge_lease (fulfillment_status, fulfillment_lease_until); diff --git a/sql/017_pending_order_guard.sql b/sql/017_pending_order_guard.sql new file mode 100644 index 0000000..5878532 --- /dev/null +++ b/sql/017_pending_order_guard.sql @@ -0,0 +1,26 @@ +-- 待支付订单约束与超时扫描优化。执行前请确认不存在重复的 unpaid/processing 订单。 + +ALTER TABLE traffic_order + MODIFY COLUMN pay_status VARCHAR(16) NOT NULL DEFAULT 'unpaid' COMMENT 'unpaid/processing/paid/cancelled/closed', + ADD KEY idx_traffic_order_pending_expire (pay_status, created_at, id); + +ALTER TABLE sim_recharge_order + ADD COLUMN closed_at DATETIME NULL COMMENT '关闭时间' AFTER paid_at, + MODIFY COLUMN payment_status VARCHAR(16) NOT NULL DEFAULT 'unpaid' COMMENT 'unpaid/processing/paid/cancelled/closed', + ADD KEY idx_sim_recharge_pending_expire (payment_status, created_at, id); + +ALTER TABLE payment_transaction + MODIFY COLUMN status VARCHAR(16) NOT NULL DEFAULT 'unpaid' COMMENT 'unpaid/processing/paid/failed/cancelled/closed', + ADD KEY idx_payment_business_status (business_type, business_order_id, status); + +ALTER TABLE traffic_order + ADD COLUMN pending_user_id BIGINT GENERATED ALWAYS AS ( + CASE WHEN pay_status IN ('unpaid', 'processing') THEN user_id ELSE NULL END + ) STORED, + ADD UNIQUE KEY uk_traffic_order_pending_user (pending_user_id); + +ALTER TABLE sim_recharge_order + ADD COLUMN pending_sim_card_id BIGINT GENERATED ALWAYS AS ( + CASE WHEN payment_status IN ('unpaid', 'processing') THEN sim_card_id ELSE NULL END + ) STORED, + ADD UNIQUE KEY uk_sim_recharge_pending_card (pending_sim_card_id); diff --git a/sql/018_sim_recharge_single_submission.sql b/sql/018_sim_recharge_single_submission.sql new file mode 100644 index 0000000..a04f15c --- /dev/null +++ b/sql/018_sim_recharge_single_submission.sql @@ -0,0 +1,15 @@ +-- SIM 续费改为支付成功后仅提交一次;异常订单转人工处理,禁止自动重试。 + +ALTER TABLE sim_recharge_order + ADD COLUMN fulfillment_started_at DATETIME NULL COMMENT '首次提交 SIMBOSS 时间' AFTER fulfillment_status, + ADD COLUMN fulfillment_message VARCHAR(128) NOT NULL DEFAULT '' COMMENT '用户可见履约提示' AFTER last_error_message, + MODIFY COLUMN fulfillment_status VARCHAR(16) NOT NULL DEFAULT 'pending' COMMENT 'pending/submitting/success/manual_review'; + +UPDATE sim_recharge_order +SET fulfillment_status = 'manual_review', + fulfillment_message = '续费处理异常,请联系管理员。', + next_attempt_at = NULL, + fulfillment_lease_token = NULL, + fulfillment_lease_until = NULL, + last_error_code = CASE WHEN last_error_code = '' THEN 'legacy_manual_review' ELSE last_error_code END +WHERE payment_status = 'paid' AND fulfillment_status <> 'success'; \ No newline at end of file diff --git a/vo/account_vo.go b/vo/account_vo.go index c2171bc..404fced 100644 --- a/vo/account_vo.go +++ b/vo/account_vo.go @@ -48,6 +48,35 @@ type TrafficOrderCreateReq struct { PackageID int64 `json:"packageId" validate:"omitempty,gt=0"` } +type PaymentCreateReq struct { + BusinessType string `json:"businessType" validate:"required,oneof=traffic_order sim_recharge_order"` + BusinessOrderID int64 `json:"businessOrderId" validate:"required,gt=0"` +} + +type PaymentCreateVO struct { + TransactionID int64 `json:"transactionId"` + BusinessType string `json:"businessType"` + BusinessOrderID int64 `json:"businessOrderId"` + Provider string `json:"provider"` + MerchantOrderNo string `json:"merchantOrderNo"` + TotalFeeFen int64 `json:"totalFeeFen"` + Status string `json:"status"` + ProviderPayload string `json:"providerPayload,omitempty"` + CallbackEnabled bool `json:"callbackEnabled"` +} + +type PaymentTransactionVO struct { + ID int64 `json:"id"` + BusinessType string `json:"businessType"` + BusinessOrderID int64 `json:"businessOrderId"` + Provider string `json:"provider"` + MerchantOrderNo string `json:"merchantOrderNo"` + TotalFeeFen int64 `json:"totalFeeFen"` + Status string `json:"status"` + CreatedAt time.Time `json:"createdAt"` + PaidAt *time.Time `json:"paidAt"` +} + // SimCardPageReq SIM 卡分页 type SimCardPageReq struct { common.Pagination @@ -61,7 +90,7 @@ type SimRechargeLogPageReq struct { RechargeStatus string `json:"rechargeStatus" form:"rechargeStatus"` } -// SimRechargeReq SIM 卡充值 +// SimRechargeReq SIM 卡续费订单 type SimRechargeReq struct { - AmountGb int `json:"amountGb" validate:"required,min=1"` + PackageID int64 `json:"packageId" validate:"required,gt=0"` } diff --git a/vo/platform_billing_vo.go b/vo/platform_billing_vo.go index 16a2d30..d646923 100644 --- a/vo/platform_billing_vo.go +++ b/vo/platform_billing_vo.go @@ -46,6 +46,16 @@ type TrafficPackagePageReq struct { Status string `json:"status" form:"status"` } +type TrafficPackageVO struct { + ID int64 `json:"id"` + Code string `json:"code"` + Name string `json:"name"` + AmountBytes int64 `json:"amountBytes"` + Price float64 `json:"price"` + ValidityDays int `json:"validityDays"` + Sort int `json:"sort"` +} + type TrafficLedgerPageReq struct { common.Pagination AccountType string `json:"accountType" form:"accountType"`