feat: package limit and notify

This commit is contained in:
wangfuduo 2026-07-16 10:30:57 +08:00
parent 968d762d54
commit 10015eac1b
10 changed files with 358 additions and 17 deletions

View File

@ -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}))
}

28
app/api/service/pkg.go Normal file
View File

@ -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
}

View File

@ -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. 如果都不可用,则提示联系客服

View File

@ -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 {

View File

@ -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
}

View File

@ -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())
}
}
}

216
app/worker/pkg.go Normal file
View File

@ -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
}

View File

@ -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()
}

View File

@ -18,4 +18,5 @@ var (
ErrExist = errors.T(11002, "数据已存在")
ErrPkgExpire = errors.T(30004, "包已过期")
ErrPkgLimit = errors.T(30005, "当前版本使用人数已满,请更新至最新版本")
)

View File

@ -8,6 +8,7 @@ var KeysDefault = mergeDefaults(
vipPayDefaults,
pkgUpdateDefaults,
pkgInstallDefaults,
pkgReportDefaults,
)
func mergeDefaults(maps ...map[string]any) map[string]any {
@ -72,6 +73,9 @@ const (
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{
@ -80,4 +84,17 @@ var pkgInstallDefaults = map[string]any{
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",
}