package service import ( "bytes" "encoding/json" "errors" "fmt" "io" "log" "net/http" "os" "path/filepath" "strings" "time" "affiliate_dash/internal/model" "affiliate_dash/internal/pkg/timeutil" "github.com/google/uuid" "gorm.io/gorm" "gorm.io/gorm/clause" ) // 充值约束:最小 10 元;1 元 = 100 积分,即 1 分人民币 = 1 积分。 const ( RechargeMinAmountCNYCents = 1000 // 10 元(单位:分) RechargePointsPerCNYCent = 1 // 每 1 分人民币兑换积分 ) type RechargeService struct { db *gorm.DB fulfill *FulfillmentService uploadDir string } func NewRechargeService(db *gorm.DB, fulfill *FulfillmentService, uploadDir string) *RechargeService { return &RechargeService{db: db, fulfill: fulfill, uploadDir: uploadDir} } type CreateRechargeInput struct { MerchantID uint ActorUserID uint AmountCNY int64 // 人民币,单位:分 Vouchers []string Note string } func (s *RechargeService) CreateRecharge(in CreateRechargeInput) (*model.RechargeApplication, error) { if in.MerchantID == 0 { return nil, errors.New("无效的商户") } if in.AmountCNY < RechargeMinAmountCNYCents { return nil, errors.New("最小充值金额为 10 元") } if len(in.Vouchers) == 0 { return nil, errors.New("请上传打款凭证截图") } cleaned := make([]string, 0, len(in.Vouchers)) for _, v := range in.Vouchers { v = strings.TrimSpace(v) if v != "" { cleaned = append(cleaned, v) } } if len(cleaned) == 0 { return nil, errors.New("请上传打款凭证截图") } if len(cleaned) > 5 { return nil, errors.New("凭证截图最多 5 张") } if len(in.Note) > 512 { return nil, errors.New("备注最长 512 个字符") } app := &model.RechargeApplication{ MerchantID: in.MerchantID, ApplicationNo: newRechargeApplicationNo(), AmountCNY: in.AmountCNY, PointsAmount: in.AmountCNY * RechargePointsPerCNYCent, Status: model.RechargeStatusPending, Vouchers: cleaned, Note: strings.TrimSpace(in.Note), } if err := s.db.Create(app).Error; err != nil { return nil, err } _ = writeAudit(s.db, &app.MerchantID, &in.ActorUserID, nil, "recharge.create", "recharge_application", fmt.Sprint(app.ID), map[string]interface{}{"application_no": app.ApplicationNo, "amount_cny": app.AmountCNY}) return app, nil } func (s *RechargeService) ListRechargeApplications(merchantID uint, page, size int, status string) ([]model.RechargeApplication, int64, error) { page, size = normalizePage(page, size) tx := s.db.Model(&model.RechargeApplication{}).Where("merchant_id = ?", merchantID) if status != "" { tx = tx.Where("status = ?", status) } var total int64 if err := tx.Count(&total).Error; err != nil { return nil, 0, err } var list []model.RechargeApplication err := tx.Preload("Merchant").Order("id DESC").Offset((page - 1) * size).Limit(size).Find(&list).Error return list, total, err } func (s *RechargeService) ListAllRechargeApplications(page, size int, status string) ([]model.RechargeApplication, int64, error) { page, size = normalizePage(page, size) tx := s.db.Model(&model.RechargeApplication{}) if status != "" { tx = tx.Where("status = ?", status) } var total int64 if err := tx.Count(&total).Error; err != nil { return nil, 0, err } var list []model.RechargeApplication err := tx.Preload("Merchant").Order("id DESC").Offset((page - 1) * size).Limit(size).Find(&list).Error return list, total, err } type ReviewRechargeInput struct { ApplicationID uint Approved bool ReviewNote string ActorUserID uint } // ReviewRecharge 审核充值申请;通过时在申请与入账共用一个事务内调用 AdjustWallet // (幂等键 = 申请单号),保证审核状态与钱包入账一致。 func (s *RechargeService) ReviewRecharge(in ReviewRechargeInput) (*model.RechargeApplication, error) { var out model.RechargeApplication err := s.db.Transaction(func(tx *gorm.DB) error { var app model.RechargeApplication if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&app, in.ApplicationID).Error; err != nil { return errors.New("充值申请不存在") } if app.Status != model.RechargeStatusPending { return errors.New("该申请已审核,不能重复操作") } status := model.RechargeStatusRejected if in.Approved { status = model.RechargeStatusApproved } now := time.Now() updates := map[string]interface{}{ "status": status, "review_note": strings.TrimSpace(in.ReviewNote), "reviewed_at": now, "reviewed_by": in.ActorUserID, } if err := tx.Model(&app).Updates(updates).Error; err != nil { return err } if in.Approved { wallet, err := s.fulfill.AdjustWallet(WalletAdjustInput{ MerchantID: app.MerchantID, ActorUserID: in.ActorUserID, Amount: app.PointsAmount, IdempotencyKey: app.ApplicationNo, Note: fmt.Sprintf("充值入账 %s(%s)", app.ApplicationNo, app.Note), }) if err != nil { return fmt.Errorf("积分入账失败:%w", err) } _ = wallet } out = app out.Status = status out.ReviewNote = strings.TrimSpace(in.ReviewNote) out.ReviewedAt = &now return writeAudit(tx, &app.MerchantID, &in.ActorUserID, nil, "recharge.review", "recharge_application", fmt.Sprint(app.ID), map[string]interface{}{"application_no": app.ApplicationNo, "approved": in.Approved, "points_amount": app.PointsAmount}) }) if err != nil { return nil, err } return &out, nil } type SaveAlertInput struct { Enabled bool ThresholdPoints int64 WebhookURL string } func (s *RechargeService) GetAlertConfig(merchantID uint) (*model.MerchantAlertConfig, error) { var cfg model.MerchantAlertConfig err := s.db.Where("merchant_id = ?", merchantID).First(&cfg).Error if errors.Is(err, gorm.ErrRecordNotFound) { return &model.MerchantAlertConfig{MerchantID: merchantID}, nil } if err != nil { return nil, err } return &cfg, nil } func (s *RechargeService) SaveAlertConfig(merchantID uint, in SaveAlertInput) (*model.MerchantAlertConfig, error) { if in.ThresholdPoints < 0 { return nil, errors.New("预警阈值不能为负") } if in.WebhookURL == "" || len(in.WebhookURL) > 1024 { return nil, errors.New("请填写 Webhook URL") } if in.Enabled { if !strings.HasPrefix(in.WebhookURL, "http://") && !strings.HasPrefix(in.WebhookURL, "https://") { return nil, errors.New("Webhook URL 必须以 http(s):// 开头") } } var cfg model.MerchantAlertConfig err := s.db.Where("merchant_id = ?", merchantID).First(&cfg).Error switch { case errors.Is(err, gorm.ErrRecordNotFound): cfg = model.MerchantAlertConfig{MerchantID: merchantID, Enabled: in.Enabled, ThresholdPoints: in.ThresholdPoints, WebhookURL: strings.TrimSpace(in.WebhookURL)} if err := s.db.Create(&cfg).Error; err != nil { return nil, err } case err != nil: return nil, err default: cfg.Enabled = in.Enabled cfg.ThresholdPoints = in.ThresholdPoints cfg.WebhookURL = strings.TrimSpace(in.WebhookURL) if err := s.db.Save(&cfg).Error; err != nil { return nil, err } } return &cfg, nil } // CheckLowBalanceAndNotify 余额低于阈值时向配置的 Webhook 推送告警。 func (s *RechargeService) CheckLowBalanceAndNotify(merchantID uint) { cfg, err := s.GetAlertConfig(merchantID) if err != nil || !cfg.Enabled || cfg.ThresholdPoints <= 0 { return } wallet, err := s.fulfill.GetWallet(merchantID) if err != nil { return } if wallet.AvailableBalance >= cfg.ThresholdPoints { return } var merchant model.Merchant if err := s.db.Select("code", "name").First(&merchant, merchantID).Error; err != nil { merchant.Name = fmt.Sprintf("商户#%d", merchantID) } payload := map[string]interface{}{ "msgtype": "text", "text": map[string]string{ "content": fmt.Sprintf("余额告警:商户 %s(#%d)当前积分余额 %d,已低于阈值 %d,请及时处理充值。", merchant.Name, merchantID, wallet.AvailableBalance, cfg.ThresholdPoints), }, } raw, _ := json.Marshal(payload) client := &http.Client{Timeout: 10 * time.Second} resp, err := client.Post(cfg.WebhookURL, "application/json", bytes.NewReader(raw)) if err != nil { log.Printf("低余额告警推送失败 merchant=%d: %v", merchantID, err) return } defer resp.Body.Close() if resp.StatusCode < 200 || resp.StatusCode >= 300 { log.Printf("低余额告警推送非 2xx merchant=%d status=%d", merchantID, resp.StatusCode) } } // SaveUploadFile 保存上传的图片到 uploadDir,返回可通过 /uploads 访问的 URL。 // 通过文件魔数校验真实图片格式(jpg/png/gif/webp),存储扩展名以实际格式为准。 func (s *RechargeService) SaveUploadFile(r io.Reader, _ string, maxBytes int64) (string, error) { data, err := io.ReadAll(io.LimitReader(r, maxBytes+1)) if err != nil { return "", errors.New("读取上传文件失败") } if len(data) == 0 { return "", errors.New("上传文件为空") } if int64(len(data)) > maxBytes { return "", fmt.Errorf("图片大小不能超过 %dMB", maxBytes/1024/1024) } ext := detectImageExt(data) if ext == "" { return "", errors.New("仅支持 jpg/png/gif/webp 图片") } name := uuid.NewString() + "." + ext if err := os.MkdirAll(s.uploadDir, 0o755); err != nil { return "", errors.New("创建上传目录失败") } if err := os.WriteFile(filepath.Join(s.uploadDir, name), data, 0o644); err != nil { return "", errors.New("保存上传文件失败") } return "/uploads/" + name, nil } // detectImageExt 根据文件头魔数识别图片格式;非法文件返回空串。 func detectImageExt(data []byte) string { switch { case len(data) >= 3 && data[0] == 0xFF && data[1] == 0xD8 && data[2] == 0xFF: return "jpg" case len(data) >= 8 && data[0] == 0x89 && data[1] == 'P' && data[2] == 'N' && data[3] == 'G': return "png" case len(data) >= 6 && data[0] == 'G' && data[1] == 'I' && data[2] == 'F' && data[3] == '8': return "gif" case len(data) >= 12 && data[0] == 'R' && data[1] == 'I' && data[2] == 'F' && data[3] == 'F' && data[8] == 'W' && data[9] == 'E' && data[10] == 'B' && data[11] == 'P': return "webp" default: return "" } } func newRechargeApplicationNo() string { return "RC" + timeutil.Now().Format(timeutil.OrderNoLayout) + strings.ReplaceAll(uuid.NewString()[:8], "-", "") }