217 lines
6.6 KiB
Go
217 lines
6.6 KiB
Go
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
|
||
}
|