From 10015eac1ba3c0a03b82029a32ed19b84886b2c3 Mon Sep 17 00:00:00 2001 From: wangfuduo Date: Thu, 16 Jul 2026 10:30:57 +0800 Subject: [PATCH] feat: package limit and notify --- app/api/controller/pkg.go | 7 +- app/api/service/pkg.go | 28 +++++ app/api/service/user.go | 14 ++- app/common/dao/pkg.go | 23 ++++ app/common/dao/user_pkg.go | 26 ++++ app/common/middle/user_pkg.go | 11 +- app/worker/pkg.go | 216 ++++++++++++++++++++++++++++++++++ app/worker/worker.go | 12 ++ base/config/error.go | 1 + base/config/keys.go | 37 ++++-- 10 files changed, 358 insertions(+), 17 deletions(-) create mode 100644 app/api/service/pkg.go create mode 100644 app/worker/pkg.go diff --git a/app/api/controller/pkg.go b/app/api/controller/pkg.go index f49fdef..44105cd 100644 --- a/app/api/controller/pkg.go +++ b/app/api/controller/pkg.go @@ -3,15 +3,20 @@ package controller import ( "net/http" "twin-api/app/api/service" + "twin-api/base/config" "git.u8t.cn/open/go-server/session" "github.com/gin-gonic/gin" + "github.com/spf13/cast" ) type Pkg struct{} // Latest 返回最新可用包的下载链接 func (p *Pkg) Latest(ctx *gin.Context) { - link := service.NewPkg().Latest() + sess := ctx.Keys[session.ContextSession].(*session.ApiSession) + limit := cast.ToInt64(sess.GetExperiment().GetParam(config.ServicePackageLimit, config.KeysDefault[config.ServicePackageLimit])) + + link := service.NewPkg().Latest(limit) ctx.JSON(http.StatusOK, session.NewRsp(gin.H{"link": link})) } diff --git a/app/api/service/pkg.go b/app/api/service/pkg.go new file mode 100644 index 0000000..e91ff72 --- /dev/null +++ b/app/api/service/pkg.go @@ -0,0 +1,28 @@ +package service + +import ( + "time" + "twin-api/app/common/dao" + "twin-api/base/config" +) + +type Pkg struct{} + +func NewPkg() *Pkg { + return &Pkg{} +} + +// Latest 按 id 倒序取最新的、未达人数上限的可用包,过期则报错,否则返回下载链接 +func (p *Pkg) Latest(limit int64) string { + pkg, err := dao.NewPackage().GetAvailableByName("", limit) + if err != nil { + panic(config.ErrDb.New().Append(err)) + } + if pkg == nil { + panic(config.ErrNoData.New().Append("package not found")) + } + if pkg.ExpireTime < time.Now().Unix() { + panic(config.ErrPkgExpire.New()) + } + return pkg.Link +} diff --git a/app/api/service/user.go b/app/api/service/user.go index 8fbdd76..bcc831d 100644 --- a/app/api/service/user.go +++ b/app/api/service/user.go @@ -23,6 +23,10 @@ func (p Params) Int(key string) int { return cast.ToInt(p[key]) } +func (p Params) Int64(key string) int64 { + return cast.ToInt64(p[key]) +} + type User struct { sess *session.ApiSession params Params @@ -124,7 +128,7 @@ func (u *User) PkgCheck() *response.Response { } if userPkg == nil { - pkg, err = dao.NewPackage().GetByPkgName("") + pkg, err = dao.NewPackage().GetAvailableByName("", u.params.Int64(config.ServicePackageLimit)) if err != nil { panic(config.ErrDb.New().Append(err)) } @@ -162,14 +166,14 @@ func (u *User) expireCheck(userPkg *model.UserPkgInfo, day int) *response.Respon if time.Unix(userPkg.Pkg.ExpireTime, 0).Before(expire) { // 提示强制更新 // 查询当前有无最新可用的包 - // 1. 先查询当前包有无可用 - if pkg, err = dao.NewPackage().GetByPkgName(userPkg.Pkg.Name); err != nil { + // 1. 先查询当前包有无可用(未满)版本 + if pkg, err = dao.NewPackage().GetAvailableByName(userPkg.Pkg.Name, u.params.Int64(config.ServicePackageLimit)); err != nil { panic(config.ErrDb.New().Append(err)) } - // 2. 无则查询其他可用包 + // 2. 无则查询其他可用(未满)包 if pkg == nil || time.Unix(pkg.ExpireTime, 0).Before(expire) { - if pkg, err = dao.NewPackage().GetByPkgName(""); err != nil { + if pkg, err = dao.NewPackage().GetAvailableByName("", u.params.Int64(config.ServicePackageLimit)); err != nil { panic(config.ErrDb.New().Append(err)) } // 3. 如果都不可用,则提示联系客服 diff --git a/app/common/dao/pkg.go b/app/common/dao/pkg.go index bf09d97..bb46688 100644 --- a/app/common/dao/pkg.go +++ b/app/common/dao/pkg.go @@ -83,6 +83,29 @@ func (p *Pkg) GetByPkgName(name string) (*model.Pkg, error) { return &res, nil } +// GetAvailableByName 按 id 倒序取最新的、未达人数上限的可用包。 +// name 为空则不限包名;limit <= 0 表示不做人数限制(等价于 GetByPkgName)。 +func (p *Pkg) GetAvailableByName(name string, limit int64) (*model.Pkg, error) { + var res model.Pkg + tx := p.db.Table(p.TableName()).Where("status = ?", model.PkgNormal) + if name != "" { + tx = tx.Where("name = ?", name) + } + if limit > 0 { + // 排除已达人数上限的包(小版本) + sub := fmt.Sprintf("(SELECT COUNT(*) FROM %s WHERE package_id = %s.id) < ?", NewUserPkg().TableName(), p.TableName()) + tx = tx.Where(sub, limit) + } + tx = tx.Order("id desc") + if err := tx.First(&res).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return nil, err + } + return &res, nil +} + func (p *Pkg) GetLatestByName(name string) (*model.Pkg, error) { var list []*model.Pkg if err := p.db.Table(p.TableName()).Where("name = ? AND status = ?", name, model.PkgNormal).Find(&list).Error; err != nil { diff --git a/app/common/dao/user_pkg.go b/app/common/dao/user_pkg.go index 1adeb3b..46bd7c9 100644 --- a/app/common/dao/user_pkg.go +++ b/app/common/dao/user_pkg.go @@ -27,6 +27,32 @@ func (p *UserPkg) Create(m *model.UserPkg) error { return p.db.Table(p.TableName()).Create(m).Error } +// CreateWithinLimit 在不超过 limit 的前提下为用户落库包记录。 +// 仅当该包(小版本)当前人数 < limit 时才创建;返回 true 表示成功,false 表示已满。 +// limit <= 0 表示不限制。 +func (p *UserPkg) CreateWithinLimit(m *model.UserPkg, limit int64) (bool, error) { + if limit > 0 { + count, err := p.CountByPackageId(m.PackageId) + if err != nil { + return false, err + } + if count >= limit { + return false, nil + } + } + if err := p.Create(m); err != nil { + return false, err + } + return true, nil +} + +// CountByPackageId 统计某个包(小版本)当前的用户数 +func (p *UserPkg) CountByPackageId(packageId int64) (int64, error) { + var count int64 + err := p.db.Table(p.TableName()).Where("package_id = ?", packageId).Count(&count).Error + return count, err +} + func (p *UserPkg) Update(m *model.UserPkg) error { return p.db.Table(p.TableName()).Save(m).Error } diff --git a/app/common/middle/user_pkg.go b/app/common/middle/user_pkg.go index 72598a9..f20231b 100644 --- a/app/common/middle/user_pkg.go +++ b/app/common/middle/user_pkg.go @@ -10,6 +10,7 @@ import ( "git.u8t.cn/open/go-server/session" "github.com/gin-gonic/gin" "github.com/smbrave/goutil" + "github.com/spf13/cast" ) func UserPkg(ctx *gin.Context) { @@ -51,8 +52,16 @@ func UserPkg(ctx *gin.Context) { userPkg.PackageId = pkg.Id userPkg.CreateTime = time.Now().Unix() userPkg.Extra = goutil.EncodeJSON(pkg) - if err = dao.NewUserPkg().Create(userPkg); err != nil { + + // 人数限制:先计数再落库,满员则拒绝该版本(引导用户更新)。 + // 非严格原子,高并发下存在极小超限窗口,业务可接受。 + limit := cast.ToInt64(sess.GetExperiment().GetParam(config.ServicePackageLimit, config.KeysDefault[config.ServicePackageLimit])) + ok, err := dao.NewUserPkg().CreateWithinLimit(userPkg, limit) + if err != nil { panic(config.ErrDb.New().Append("user_pkg: create error: ", err)) } + if !ok { + panic(config.ErrPkgLimit.New()) + } } } diff --git a/app/worker/pkg.go b/app/worker/pkg.go new file mode 100644 index 0000000..db97977 --- /dev/null +++ b/app/worker/pkg.go @@ -0,0 +1,216 @@ +package worker + +import ( + "fmt" + "math" + "strings" + "time" + "twin-api/app/common/dao" + "twin-api/app/common/model" + "twin-api/base/config" + "twin-api/base/global" + + "git.u8t.cn/open/go-server/session" + gsUtils "git.u8t.cn/open/go-server/utils" + "github.com/go-redis/redis" + log "github.com/sirupsen/logrus" + "github.com/spf13/cast" +) + +const ( + pkgLimitNotifyKey = "twin:pkg_limit:notify:%d" // 退避状态,value 形如 "lastSent:level" + pkgLimitNotifyBase = int64(3600) // 退避基准间隔:1 小时(秒) + pkgLimitNotifyMax = int64(24 * 3600) // 退避间隔上限:24 小时(秒) + pkgLimitNotifyMaxTimes = 5 // 单个包版本累计最多发送次数 + pkgLimitNotifyExpire = 30 * 24 * time.Hour // 状态自动过期,兜底清理 +) + +type PkgLimit struct{} + +// CheckLimit 巡检每个包(小版本)的使用人数,达到 limit*ratio 阈值时通知配置的接收人 +func (p *PkgLimit) CheckLimit() { + exp := session.NewExperiment(nil, nil) + limit := cast.ToInt64(exp.GetParam(config.ServicePackageLimit, config.KeysDefault[config.ServicePackageLimit])) + ratio := cast.ToFloat64(exp.GetParam(config.ServicePackageLimitRatio, config.KeysDefault[config.ServicePackageLimitRatio])) + receivers := parseReceivers(exp.GetParam(config.ServicePackageLimitReceiver, config.KeysDefault[config.ServicePackageLimitReceiver])) + + threshold := int64(float64(limit) * ratio) + if limit <= 0 || threshold <= 0 || len(receivers) == 0 { + return + } + + pkgs, _, err := dao.NewPackage().Query("", "", model.PkgNormal, 1, math.MaxInt16, "") + if err != nil { + log.Errorf("pkg_limit: query packages err:%v", err) + return + } + if len(pkgs) == 0 { + return + } + + ids := make([]int64, 0, len(pkgs)) + for _, pkg := range pkgs { + ids = append(ids, pkg.Id) + } + + counts, err := dao.NewUserPkg().CountByPackageIds(ids) + if err != nil { + log.Errorf("pkg_limit: count users err:%v", err) + return + } + + now := time.Now().Unix() + for _, pkg := range pkgs { + count := counts[pkg.Id] + if count < threshold { + // 掉回阈值以下则清空退避状态,下次越线重新从头计时 + p.resetNotify(pkg.Id) + continue + } + seq, ok := p.allowNotify(pkg.Id, now) + if !ok { + continue + } + p.notify(receivers, pkg, count, limit, seq) + } +} + +func (p *PkgLimit) notify(receivers []string, pkg *model.Pkg, count, limit int64, seq int) { + // 指纹带上发送序号,保证每次退避通知都能真正下发(否则会被消息中台按指纹去重丢弃) + fingerprint := fmt.Sprintf("pkg_limit_%d_%d", pkg.Id, seq) + msg := strings.Join([]string{ + "【包版本使用人数预警】", + fmt.Sprintf("包名:%s", pkg.Name), + fmt.Sprintf("版本:%s", pkg.Version), + fmt.Sprintf("当前人数:%d / %d", count, limit), + }, "\n") + + if err := global.GetMessage().Send(strings.Join(receivers, ","), msg, fingerprint); err != nil { + log.Errorf("pkg_limit: send message err pkg=%d :%v", pkg.Id, err) + } +} + +// allowNotify 基于指数退避判断当前是否应发送预警,并返回本次发送序号。 +// 首次越线立即发送;之后间隔按 1h、2h、4h…翻倍,上限 24h;累计最多发送 5 次。 +func (p *PkgLimit) allowNotify(pkgId, now int64) (int, bool) { + rdb := global.GetRedis() + key := fmt.Sprintf(pkgLimitNotifyKey, pkgId) + + val, err := rdb.Get(key).Result() + if err == redis.Nil { + p.saveNotify(key, now, 1) + return 1, true + } + if err != nil { + // 读状态失败时不发送,避免退化成每小时轰炸 + log.Errorf("pkg_limit: read notify state err pkg=%d :%v", pkgId, err) + return 0, false + } + + lastSent, level := parseNotify(val) + if level >= pkgLimitNotifyMaxTimes { + // 已达最大发送次数,停止退避通知 + return 0, false + } + if now-lastSent < backoffInterval(level) { + return 0, false + } + level++ + p.saveNotify(key, now, level) + return level, true +} + +func (p *PkgLimit) saveNotify(key string, now int64, level int) { + if err := global.GetRedis().Set(key, fmt.Sprintf("%d:%d", now, level), pkgLimitNotifyExpire).Err(); err != nil { + log.Errorf("pkg_limit: save notify state err %s :%v", key, err) + } +} + +func (p *PkgLimit) resetNotify(pkgId int64) { + if err := global.GetRedis().Del(fmt.Sprintf(pkgLimitNotifyKey, pkgId)).Err(); err != nil { + log.Errorf("pkg_limit: reset notify state err pkg=%d :%v", pkgId, err) + } +} + +func parseNotify(val string) (lastSent int64, level int) { + if parts := strings.SplitN(val, ":", 2); len(parts) == 2 { + return cast.ToInt64(parts[0]), cast.ToInt(parts[1]) + } + return 0, 1 +} + +// backoffInterval 第 level 次发送后需等待的间隔(秒):base * 2^(level-1),上限 max +func backoffInterval(level int) int64 { + if level < 1 { + level = 1 + } + shift := level - 1 + if shift >= 31 { + return pkgLimitNotifyMax + } + interval := pkgLimitNotifyBase << uint(shift) + if interval <= 0 || interval > pkgLimitNotifyMax { + return pkgLimitNotifyMax + } + return interval +} + +// DailyReport 每日汇总所有可用包及各版本的过期情况,发送报告给配置的接收人 +func (p *PkgLimit) DailyReport() { + exp := session.NewExperiment(nil, nil) + receivers := parseReceivers(exp.GetParam(config.ServicePackageReportReceiver, config.KeysDefault[config.ServicePackageReportReceiver])) + if len(receivers) == 0 { + return + } + + pkgs, _, err := dao.NewPackage().Query("", "", model.PkgNormal, 1, math.MaxInt16, "name asc, id asc") + if err != nil { + log.Errorf("pkg_report: query packages err:%v", err) + return + } + if len(pkgs) == 0 { + return + } + + now := time.Now().Unix() + lines := []string{ + "【可用包过期报告】", + fmt.Sprintf("统计时间:%s", gsUtils.TimeToDateTime(now)), + "", + } + for _, pkg := range pkgs { + remain := "已过期" + if pkg.ExpireTime > now { + remain = fmt.Sprintf("%d 天", (pkg.ExpireTime-now)/86400) + } + lines = append(lines, fmt.Sprintf( + "包名:%s | 版本:%s | 创建时间:%s | 过期时间:%s | 距过期:%s", + pkg.Name, pkg.Version, + gsUtils.TimeToDateTime(pkg.CreateTime), + gsUtils.TimeToDateTime(pkg.ExpireTime), + remain, + )) + } + + // 按天生成指纹,保证同一天只发一封、避免重复 + fingerprint := fmt.Sprintf("pkg_report_%s", time.Now().Format("20060102")) + if err := global.GetMessage().Send(strings.Join(receivers, ","), strings.Join(lines, "\n"), fingerprint); err != nil { + log.Errorf("pkg_report: send message err:%v", err) + } +} + +// parseReceivers 兼容「逗号分隔字符串」与「字符串数组」两种配置形式 +func parseReceivers(v any) []string { + receivers := cast.ToStringSlice(v) + if len(receivers) == 1 { + receivers = strings.Split(receivers[0], ",") + } + + out := make([]string, 0, len(receivers)) + for _, r := range receivers { + if r = strings.TrimSpace(r); r != "" { + out = append(out, r) + } + } + return out +} diff --git a/app/worker/worker.go b/app/worker/worker.go index 61f3f6a..87b153b 100644 --- a/app/worker/worker.go +++ b/app/worker/worker.go @@ -14,5 +14,17 @@ func NewWorker() *Worker { func (w *Worker) Run(cron *gocron.Scheduler) { Gs.InitWorker(cron) + + // 每小时巡检各包版本使用人数,达到阈值时预警通知 + var pkgLimit PkgLimit + cron.Every(1).Hour().Do(func() { + pkgLimit.CheckLimit() + }) + + // 每天 18:00 汇总所有可用包各版本的过期情况,发送报告 + cron.Every(1).Day().At("18:00").Do(func() { + pkgLimit.DailyReport() + }) + cron.StartAsync() } diff --git a/base/config/error.go b/base/config/error.go index 45bd470..ef7f4f6 100644 --- a/base/config/error.go +++ b/base/config/error.go @@ -18,4 +18,5 @@ var ( ErrExist = errors.T(11002, "数据已存在") ErrPkgExpire = errors.T(30004, "包已过期") + ErrPkgLimit = errors.T(30005, "当前版本使用人数已满,请更新至最新版本") ) diff --git a/base/config/keys.go b/base/config/keys.go index 5492324..8ad56b5 100644 --- a/base/config/keys.go +++ b/base/config/keys.go @@ -8,6 +8,7 @@ var KeysDefault = mergeDefaults( vipPayDefaults, pkgUpdateDefaults, pkgInstallDefaults, + pkgReportDefaults, ) func mergeDefaults(maps ...map[string]any) map[string]any { @@ -67,17 +68,33 @@ var pkgUpdateDefaults = map[string]any{ // ==================== 包缺失与安装 ==================== const ( - ServiceUserPackageLack = "server.user.package.lack" // 未查询到用户包信息提示语 - ServicePackageLackButton = "server.package.lack.button" // 用户包缺失按钮文字 - ServicePackageLackUrl = "server.package.lack.url" // 用户包缺失下载地址 - ServiceContactCustomer = "server.contact.customer" // 联系客服提示语 - ServiceInstallButton = "server.install.button" // 安装按钮文字 + ServiceUserPackageLack = "server.user.package.lack" // 未查询到用户包信息提示语 + ServicePackageLackButton = "server.package.lack.button" // 用户包缺失按钮文字 + ServicePackageLackUrl = "server.package.lack.url" // 用户包缺失下载地址 + ServiceContactCustomer = "server.contact.customer" // 联系客服提示语 + ServiceInstallButton = "server.install.button" // 安装按钮文字 + ServicePackageLimit = "server.package.limit" // 每个包的每个小版本人数限制 + ServicePackageLimitRatio = "server.package.limit.ratio" // 人数达到 limit*ratio 时触发预警通知(0~1) + ServicePackageLimitReceiver = "server.package.limit.receiver" // 人数预警通知接收人(逗号分隔) ) var pkgInstallDefaults = map[string]any{ - ServiceUserPackageLack: "请您重新安装", - ServicePackageLackButton: "安装", - ServicePackageLackUrl: "", - ServiceContactCustomer: "安装包缺失,请联系客服解决:wangfuduo@batiao8.com", - ServiceInstallButton: "安装", + ServiceUserPackageLack: "请您重新安装", + ServicePackageLackButton: "安装", + ServicePackageLackUrl: "", + ServiceContactCustomer: "安装包缺失,请联系客服解决:wangfuduo@batiao8.com", + ServiceInstallButton: "安装", + ServicePackageLimit: "10000", + ServicePackageLimitRatio: "0.8", + ServicePackageLimitReceiver: "wangfuduo,pengguangjian", +} + +// ==================== 包过期报告 ==================== + +const ( + ServicePackageReportReceiver = "server.package.report.receiver" // 每日包过期报告接收人(逗号分隔) +) + +var pkgReportDefaults = map[string]any{ + ServicePackageReportReceiver: "wangfuduo,pengguangjian", }