diff --git a/apps/services/comm/appscope.go b/apps/services/comm/appscope.go new file mode 100644 index 00000000..3e106d64 --- /dev/null +++ b/apps/services/comm/appscope.go @@ -0,0 +1,33 @@ +package comm + +import ( + "os" + "strings" +) + +// 部署身份(应用名/区域)解析。 +// +// 同一份部署里历史上存在两个表达「我是哪个应用」的环境变量: +// - APP_NAME —— 服务配置托管 / 会议记录编排 / 渠道分发等按应用作用域读配置时用 +// - ANALYZE_APP_NAME —— analyze 模块上报统计快照时用(modules.analyze.AppName) +// +// 两者语义完全相同(都要求「与 console 注册表 app_registry 的 app_name 一致」), +// 但线上部署普遍只填了 ANALYZE_APP_NAME、把 APP_NAME 留空——于是所有按作用域读的接口 +// (渠道分发 user_getchannelapp(s)、user_getagents_v2、user_getthirdsvcs_v2)都退化成 +// 只读「全局默认」行,后台明明按应用配好了却下发不到客户端,且不报错、极难排查。 +// +// 这里统一收口:APP_NAME 优先,留空则回退 ANALYZE_APP_NAME;都没有才是真正的全局默认("")。 + +// AppName 返回本部署的应用名(console 注册表 app_registry.app_name)。 +// 返回 "" 表示未声明部署身份,调用方按「全局默认」作用域处理。 +func AppName() string { + if v := strings.TrimSpace(os.Getenv("APP_NAME")); v != "" { + return v + } + return strings.TrimSpace(os.Getenv("ANALYZE_APP_NAME")) +} + +// AppRegion 返回本部署的区域短代码(cn/us/ea/sea…),未配置返回 ""。 +func AppRegion() string { + return strings.TrimSpace(os.Getenv("APP_REGION")) +} diff --git a/apps/services/comm/const.go b/apps/services/comm/const.go index 5f742259..2deb7e1a 100644 --- a/apps/services/comm/const.go +++ b/apps/services/comm/const.go @@ -82,10 +82,8 @@ const ( // TableFactoryDevice = "factorydevice" //厂家设备表 TableUseRecordLog = "uselog" //用户使用统计日志 - TableAppStatDaily = "app_stat_daily" //App综合统计-日表 - TableAppStatMonthly = "app_stat_monthly" //App综合统计-月表 - TableAppStatYearly = "app_stat_yearly" //App综合统计-年表 - TableAppStatGlobal = "app_stat_global" //App综合统计-全局累计(单行 stat_date=all) + // app_stat_daily/monthly/yearly/global 四张「App 综合运营统计」表已移除:运营统计整体迁到 + // 运营后台(analyze 模块把当日快照经 NATS 推到 console 主库 stats_global_day),app 服务不再产出。 TableConsoleLog = "console_log" //后台操作日志表 TableAdminResLog = "admin_resource_log" //后台账号资源流水(超管→管理员/代理/运营) TableAllhelpTask = "allhelp_task" //用户定时提醒任务表 diff --git a/apps/services/lego/sys/mysql/core.go b/apps/services/lego/sys/mysql/core.go index 1be15b7a..aae3a405 100644 --- a/apps/services/lego/sys/mysql/core.go +++ b/apps/services/lego/sys/mysql/core.go @@ -13,6 +13,8 @@ type ( Exec(sql string, values ...interface{}) (tx *gorm.DB) CreateTable(tName string, model any) (err error) AutoIncrementStart(tName, column string, start uint64) (err error) + // AutoIncrementFloor 抬高自增主键下限(表里已有数据也生效,已有行 id 不动)。 + AutoIncrementFloor(tName, column string, start uint64) (err error) FindOne(tName string, model any, query interface{}, args ...interface{}) (err error) // FindOnePrimary 与 FindOne 相同,但强制走主库(读写分离下绕过只读副本)。 // 用于"读后写同一行"等强一致场景,避免读到复制延迟内的旧数据。MySQL 无副本时等同 FindOne。 @@ -79,6 +81,10 @@ func AutoIncrementStart(tName, column string, start uint64) (err error) { return defsys.AutoIncrementStart(tName, column, start) } +func AutoIncrementFloor(tName, column string, start uint64) (err error) { + return defsys.AutoIncrementFloor(tName, column, start) +} + func Delete(tName string, query interface{}, args ...interface{}) (err error) { return defsys.Delete(tName, query, args...) } diff --git a/apps/services/lego/sys/mysql/mysql.go b/apps/services/lego/sys/mysql/mysql.go index ceaab3ce..2708b539 100644 --- a/apps/services/lego/sys/mysql/mysql.go +++ b/apps/services/lego/sys/mysql/mysql.go @@ -34,6 +34,13 @@ func (this *MySql) AutoIncrementStart(tName, column string, start uint64) (err e return } +// AutoIncrementFloor 抬高自增主键下限:新插入的行 id >= start,已有行不动。 +// MySQL 的 ALTER TABLE ... AUTO_INCREMENT 本身就是下限语义(低于当前最大值+1 时被静默忽略), +// 所以与 AutoIncrementStart 同一条语句;两个方法并存只为与 postgres 驱动的接口对齐。 +func (this *MySql) AutoIncrementFloor(tName, column string, start uint64) (err error) { + return this.AutoIncrementStart(tName, column, start) +} + func (this *MySql) Exec(sql string, values ...interface{}) (tx *gorm.DB) { tx = this.db.Exec(sql, values...) return diff --git a/apps/services/lego/sys/postgres/core.go b/apps/services/lego/sys/postgres/core.go index 8bb3e81f..529ceaee 100644 --- a/apps/services/lego/sys/postgres/core.go +++ b/apps/services/lego/sys/postgres/core.go @@ -13,6 +13,8 @@ type ( Exec(sql string, values ...interface{}) (tx *gorm.DB) CreateTable(tName string, model any) (err error) AutoIncrementStart(tName, column string, start uint64) (err error) + // AutoIncrementFloor 抬高自增主键下限(表里已有数据也生效,已有行 id 不动)。 + AutoIncrementFloor(tName, column string, start uint64) (err error) FindOne(tName string, model any, query interface{}, args ...interface{}) (err error) // FindOnePrimary 与 FindOne 相同,但强制走主库(绕过只读副本)。 // 用于"读后写同一行"等强一致场景,避免读到复制延迟内的旧数据。未配置副本时与 FindOne 等价。 @@ -79,6 +81,10 @@ func AutoIncrementStart(tName, column string, start uint64) (err error) { return defsys.AutoIncrementStart(tName, column, start) } +func AutoIncrementFloor(tName, column string, start uint64) (err error) { + return defsys.AutoIncrementFloor(tName, column, start) +} + func Delete(tName string, query interface{}, args ...interface{}) (err error) { return defsys.Delete(tName, query, args...) } diff --git a/apps/services/lego/sys/postgres/postgres.go b/apps/services/lego/sys/postgres/postgres.go index d8c99b10..4476f988 100644 --- a/apps/services/lego/sys/postgres/postgres.go +++ b/apps/services/lego/sys/postgres/postgres.go @@ -182,6 +182,32 @@ func (this *Postgres) AutoIncrementStart(tName, column string, start uint64) (er return } +// AutoIncrementFloor 抬高表自增主键的下限:保证之后新插入的行 id >= start。 +// +// 与 AutoIncrementStart 的区别是「表里已有数据也照样生效」——已有行的 id 原样保留, +// 只把序列推到 max(当前最大 id, start-1),因此新行从 start 起(或从已有最大值 +1 起, +// 若它本就更大)。用于给已上线的表补上带业务含义的 id 段(如方案商 0xAB01 起): +// 只用 AutoIncrementStart 的话,表一旦有过一行就永远停在 1、2、3…。 +// +// 语义与 MySQL 的 `ALTER TABLE ... AUTO_INCREMENT = n` 一致(MySQL 侧本就是下限语义)。 +func (this *Postgres) AutoIncrementFloor(tName, column string, start uint64) (err error) { + if start == 0 { + return + } + // setval(..., v, true) 让下一个 nextval = v+1,故取 v = GREATEST(表内最大 id, 序列当前值, start-1)。 + // 三者都要参与: + // - 表内最大 id —— 保证不会退回去撞上已有行; + // - 序列当前值 —— 保证不会把已经发到更高位的序列往回拨(回滚过的插入会让序列超前于最大 id); + // - start-1 —— 本次要抬到的下限。 + err = this.db.Exec( + `SELECT setval(pg_get_serial_sequence(?, ?), GREATEST( + (SELECT COALESCE(MAX(`+column+`), 0) FROM `+tName+`), + COALESCE(pg_sequence_last_value(pg_get_serial_sequence(?, ?)), 0), + ?), true)`, + tName, column, tName, column, start-1).Error + return +} + // 获取表对象 func (this *Postgres) Table(tName string) (tx *gorm.DB) { return this.db.Table(tName) diff --git a/apps/services/modules/api/api_getglobalstats.go b/apps/services/modules/api/api_getglobalstats.go index 37237362..9f44e812 100644 --- a/apps/services/modules/api/api_getglobalstats.go +++ b/apps/services/modules/api/api_getglobalstats.go @@ -7,8 +7,16 @@ import ( ) // GetGlobalStats 查询全局累计统计 +// +// stat 字段恒为空:它原先读 app_stat_global,而 App 运营统计已整体迁到运营后台 +// (console 主库 stats_global_day,后台读 api_getstatssummary / api_getstatstrend), +// app 服务侧四张 app_stat_* 表连同写入链路已删除。 +// +// 累计注册/充值人数仍然真实返回——它们是对 user / payorder 两张**业务表**直接 COUNT, +// 不依赖被删掉的统计表,留着不花什么代价,也免得老后台上这两个数字凭空变 0。 +// // @Summary 查询全局累计统计 -// @Description 返回生命周期累计的 App 综合统计,以及实时计算的累计注册/充值人数 +// @Description 返回实时计算的累计注册/充值人数;stat 字段已停用(App 运营统计迁至运营后台) // @Tags API // @Accept json // @Produce json @@ -17,16 +25,10 @@ import ( // @Router /web/api/api_getglobalstats [post] func (this *apiComp) GetGlobalStats(session comm.IUserSession, req *pb.ApiGetGlobalStatsReq) (resp *pb.ApiGetGlobalStatsResp, errdata *pb.ErrorData) { var ( - stat *pb.DBAppStat totalUsers int64 totalPayUser int64 err error ) - if stat, err = this.module.model.getGlobalStat(); err != nil { - errdata = &pb.ErrorData{Code: pb.ErrorCode_DBError, Message: err.Error()} - this.module.Error("GetGlobalStats fail", log.Field{Key: "err", Value: err.Error()}) - return - } if totalUsers, err = this.module.model.countTotalUsers(); err != nil { errdata = &pb.ErrorData{Code: pb.ErrorCode_DBError, Message: err.Error()} this.module.Error("countTotalUsers fail", log.Field{Key: "err", Value: err.Error()}) @@ -38,7 +40,7 @@ func (this *apiComp) GetGlobalStats(session comm.IUserSession, req *pb.ApiGetGlo return } resp = &pb.ApiGetGlobalStatsResp{ - Stat: stat, + Stat: &pb.DBAppStat{StatDate: "all"}, // 占位空行:运营统计已迁至运营后台 TotalUserCount: totalUsers, TotalPayUserCount: totalPayUser, } diff --git a/apps/services/modules/api/api_getpaystats.go b/apps/services/modules/api/api_getpaystats.go index a8fc08a9..0d379f1b 100644 --- a/apps/services/modules/api/api_getpaystats.go +++ b/apps/services/modules/api/api_getpaystats.go @@ -2,27 +2,18 @@ package api import ( "yunyan/comm" - "yunyan/lego/sys/log" "yunyan/pb" ) -// 查询App综合统计数据(日/月/年) +// GetAppStats 查询App综合统计数据(日/月/年) +// +// 【已停用】App 运营统计整体迁到运营后台:analyze 模块把当日快照经 NATS 推到 console 主库 +// stats_global_day,后台读 console 的 api_getstatssummary / api_getstatstrend。 +// app 服务侧的 app_stat_daily/monthly/yearly/global 四张表连同写入链路已删除, +// 这里没有可查的数据源,恒定返回空列表。 +// +// 路由保留只为老调用方不至于 404;确认无人再调后可连同 pb 定义一起删掉。 func (this *apiComp) GetAppStats(session comm.IUserSession, req *pb.ApiGetAppStatsReq) (resp *pb.ApiGetAppStatsResp, errdata *pb.ErrorData) { - var ( - stats []*pb.DBAppStat - err error - ) - if req.StatType == 0 { - req.StatType = 1 // 默认日统计 - } - if stats, err = this.module.model.getAppStats(req.StatType, req.StartDate, req.EndDate); err != nil { - errdata = &pb.ErrorData{ - Code: pb.ErrorCode_DBError, - Message: err.Error(), - } - this.module.Error("GetAppStats fail", log.Field{Key: "uid", Value: session.GetUserId()}, log.Field{Key: "err", Value: err.Error()}) - return - } - resp = &pb.ApiGetAppStatsResp{Stats: stats} + resp = &pb.ApiGetAppStatsResp{Stats: make([]*pb.DBAppStat, 0)} return } diff --git a/apps/services/modules/api/api_rebuildstats.go b/apps/services/modules/api/api_rebuildstats.go index 36194075..ad2580f8 100644 --- a/apps/services/modules/api/api_rebuildstats.go +++ b/apps/services/modules/api/api_rebuildstats.go @@ -4,28 +4,25 @@ import ( "yunyan/comm" "yunyan/lego/sys/log" "yunyan/pb" - "sort" "strings" "time" ) -// RebuildStats 全量重建 App 综合统计(仅超管) +// RebuildStats 全量重建统计(仅超管) // -// 清空 app_stat_daily/monthly/yearly/global 后,从以下源表全量重算: -// - payorder:order_count(create_time)、failed_order_count(create_time, status=FAILED)、 -// pay_amount/pay_count/pay_user_count(pay_time, status=PAID) -// - user:new_user_count(createtime) -// - useruselog:add_vip_days/add_ai_integral/add_trade_second/add_meet_second(ts) -// - userstatistics(仅 global 行):ai/trade/meet 累计消耗 -// - userdevice(仅 global 行):bind_device_count +// 【已缩小范围】原先此接口还会清空并重算 app_stat_daily/monthly/yearly/global 四张 +// 「App 综合运营统计」表。运营统计已整体迁到运营后台(analyze 模块把当日快照经 NATS 推到 +// console 主库 stats_global_day,后台读 api_getstatssummary / api_getstatstrend), +// app 服务不再产出、也不再存运营统计,那部分逻辑连同四张表一并移除。 // -// 无法重建的字段(无历史时间戳/审计)将留 0: +// 现在仍然重建的是两张**业务表**,与运营统计无关、且各有在线读者: +// - userstatistics —— 用户维度累计消耗,用户排行榜依赖;从 useruselog + echomeet_record 重算 +// - product_stat —— 产品激活数/绑定用户数,代理仪表盘 api_getagentdashboard 依赖;从 userdevice 重算 // -// login_count、active_user_count、active_device_count; -// ai_chat/trade/meet 消耗的日/月/年分布(只能汇总到 global 行)。 +// 响应里的 daily/monthly/yearly 行数恒为 0——字段保留只为不破坏老调用方的结构。 // // @Summary 重建统计数据 -// @Description 清空并从源表重算 App 综合统计(仅超管) +// @Description 从源表重算 userstatistics 与 product_stat(仅超管);App 运营统计已迁至运营后台,不在此重建 // @Tags API // @Accept json // @Produce json @@ -36,7 +33,7 @@ func (this *apiComp) RebuildStats(session comm.IUserSession, req *pb.ApiRebuildS start := time.Now() resp = &pb.ApiRebuildStatsResp{} - // 0. 先从消耗日志重建 userstatistics(排行榜 + 后续 global 汇总均依赖此表) + // 1. 从消耗日志重建 userstatistics(用户排行榜依赖) uRows, uErr := this.module.model.rebuildUserStatisticsFromLogs() if uErr != nil { errdata = dbErr(uErr, "rebuildUserStatisticsFromLogs") @@ -46,170 +43,7 @@ func (this *apiComp) RebuildStats(session comm.IUserSession, req *pb.ApiRebuildS log.Field{Key: "rows", Value: uRows}, ) - // 1. 各源表分组聚合(DB 端 GROUP BY,避免拉全表到 Go) - orderCreate, err := this.module.model.aggOrderCreate() - if err != nil { - errdata = dbErr(err, "aggOrderCreate") - return - } - orderFailed, err := this.module.model.aggOrderFailed() - if err != nil { - errdata = dbErr(err, "aggOrderFailed") - return - } - orderPaid, err := this.module.model.aggOrderPaid() - if err != nil { - errdata = dbErr(err, "aggOrderPaid") - return - } - userSignup, err := this.module.model.aggUserSignup() - if err != nil { - errdata = dbErr(err, "aggUserSignup") - return - } - useLog, err := this.module.model.aggUseLogGrant() - if err != nil { - errdata = dbErr(err, "aggUseLogGrant") - return - } - - // 2. 合并到 daily map(key=YYYY-MM-DD) - daily := make(map[string]*pb.DBAppStat) - getOrInit := func(d string) *pb.DBAppStat { - row, ok := daily[d] - if !ok { - row = &pb.DBAppStat{StatDate: d} - daily[d] = row - } - return row - } - for _, r := range orderCreate { - getOrInit(r.D).OrderCount += int32(r.Count) - } - for _, r := range orderFailed { - getOrInit(r.D).FailedOrderCount += int32(r.Count) - } - for _, r := range orderPaid { - row := getOrInit(r.D) - row.PayAmount += r.Amount - row.PayCount += int32(r.Count) - row.PayUserCount += int32(r.UserCount) - } - for _, r := range userSignup { - getOrInit(r.D).NewUserCount += int32(r.Count) - } - for _, r := range useLog { - row := getOrInit(r.D) - row.AddVipDays += r.AddVipDay - row.AddAiIntegral += r.AddAgentInt - row.AddTradeSecond += r.AddTradeSec - row.AddMeetSecond += r.AddMeetSec - } - - // 3. 按日期排序,并累计 total_user_count 快照 - dailyRows := make([]*pb.DBAppStat, 0, len(daily)) - for _, r := range daily { - dailyRows = append(dailyRows, r) - } - sort.Slice(dailyRows, func(i, j int) bool { return dailyRows[i].StatDate < dailyRows[j].StatDate }) - var cumUser int64 - now := time.Now().Unix() - for _, r := range dailyRows { - cumUser += int64(r.NewUserCount) - r.TotalUserCount = cumUser - r.UpdateTime = now - } - - // 4. 月 / 年汇总(基于已修正 total_user_count 的日行;快照字段取该期末值) - monthlyRows := rollupByPrefix(dailyRows, 7) // YYYY-MM - yearlyRows := rollupByPrefix(dailyRows, 4) // YYYY - - // 5. global 行:日聚合无法覆盖的字段从生命周期 SQL 单独取 - gPayAmount, gPayCount, gPayUserCount, gOrderCount, gFailedCount, err := this.module.model.sumLifetimePayments() - if err != nil { - errdata = dbErr(err, "sumLifetimePayments") - return - } - gVipDay, gAiInt, gTradeSec, gMeetSec, err2 := this.module.model.sumLifetimeUseLog() - if err2 != nil { - errdata = dbErr(err2, "sumLifetimeUseLog") - return - } - gAiChat, gAiUp, gAiDown, gTradeCount, gTradeTime, gTradeWords, gMeetCount, gMeetTime, err3 := this.module.model.sumLifetimeUserStats() - if err3 != nil { - errdata = dbErr(err3, "sumLifetimeUserStats") - return - } - gTotalUsers, err := this.module.model.countTotalUsers() - if err != nil { - errdata = dbErr(err, "countTotalUsers") - return - } - gDevices, err := this.module.model.countUserDevices() - if err != nil { - errdata = dbErr(err, "countUserDevices") - return - } - gVipUsers, err := this.module.model.countActiveVipUsers(now) - if err != nil { - errdata = dbErr(err, "countActiveVipUsers") - return - } - - global := &pb.DBAppStat{ - StatDate: GlobalStatDate, - PayAmount: gPayAmount, - OrderCount: int32(gOrderCount), - PayCount: int32(gPayCount), - FailedOrderCount: int32(gFailedCount), - PayUserCount: int32(gPayUserCount), - AddVipDays: gVipDay, - AddAiIntegral: gAiInt, - AddTradeSecond: gTradeSec, - AddMeetSecond: gMeetSec, - AiChatCount: gAiChat, - AiUpToken: gAiUp, - AiDownToken: gAiDown, - TradeCount: gTradeCount, - TradeTime: gTradeTime, - TradeWords: gTradeWords, - MeetCount: gMeetCount, - MeetTime: gMeetTime, - NewUserCount: int32(gTotalUsers), - TotalUserCount: gTotalUsers, - VipUserCount: int32(gVipUsers), - BindDeviceCount: int32(gDevices), - UpdateTime: now, - } - - // 6. 原子改写:在 stat consumer 锁内 TRUNCATE + INSERT + ReloadCache, - // 避免消费者拿着旧的内存缓存覆盖刚重建的数据。 - werr := this.module.statConsumer.WithLock(func() error { - if e := this.module.model.truncateStatTables(); e != nil { - return e - } - if e := this.module.model.insertStatRows(comm.TableAppStatDaily, dailyRows); e != nil { - return e - } - if e := this.module.model.insertStatRows(comm.TableAppStatMonthly, monthlyRows); e != nil { - return e - } - if e := this.module.model.insertStatRows(comm.TableAppStatYearly, yearlyRows); e != nil { - return e - } - if e := this.module.model.insertStatRows(comm.TableAppStatGlobal, []*pb.DBAppStat{global}); e != nil { - return e - } - // 仍在锁内:直接刷新 consumer 内存缓存 - this.module.statConsumer.loadCache(time.Now()) - return nil - }) - if werr != nil { - errdata = dbErr(werr, "rebuild atomic") - return - } - - // 7. 重建产品激活/绑定统计:全量扫描 userdevice 落库 product_stat, + // 2. 重建产品激活/绑定统计:全量扫描 userdevice 落库 product_stat, // 供代理仪表盘(api_getagentdashboard)按产品读取激活数 / 绑定用户数。 psRows, psErr := this.module.model.rebuildProductStats() if psErr != nil { @@ -220,56 +54,15 @@ func (this *apiComp) RebuildStats(session comm.IUserSession, req *pb.ApiRebuildS log.Field{Key: "rows", Value: psRows}, ) - resp.DailyRows = int32(len(dailyRows)) - resp.MonthlyRows = int32(len(monthlyRows)) - resp.YearlyRows = int32(len(yearlyRows)) resp.CostMs = time.Since(start).Milliseconds() - this.module.Info("RebuildStats done", - log.Field{Key: "daily", Value: resp.DailyRows}, - log.Field{Key: "monthly", Value: resp.MonthlyRows}, - log.Field{Key: "yearly", Value: resp.YearlyRows}, + log.Field{Key: "userstatistics", Value: uRows}, + log.Field{Key: "product_stat", Value: psRows}, log.Field{Key: "cost_ms", Value: resp.CostMs}, ) return } -// rollupByPrefix 把 daily 行按日期前缀 (7=YYYY-MM, 4=YYYY) 聚合成月/年行。 -// 快照字段 total_user_count 取该期内最大值(= 期末累计用户数)。 -func rollupByPrefix(daily []*pb.DBAppStat, prefixLen int) []*pb.DBAppStat { - bucket := make(map[string]*pb.DBAppStat) - for _, d := range daily { - if len(d.StatDate) < prefixLen { - continue - } - key := d.StatDate[:prefixLen] - row, ok := bucket[key] - if !ok { - row = &pb.DBAppStat{StatDate: key, UpdateTime: d.UpdateTime} - bucket[key] = row - } - row.PayAmount += d.PayAmount - row.OrderCount += d.OrderCount - row.PayCount += d.PayCount - row.FailedOrderCount += d.FailedOrderCount - row.PayUserCount += d.PayUserCount - row.AddVipDays += d.AddVipDays - row.AddAiIntegral += d.AddAiIntegral - row.AddTradeSecond += d.AddTradeSecond - row.AddMeetSecond += d.AddMeetSecond - row.NewUserCount += d.NewUserCount - if d.TotalUserCount > row.TotalUserCount { - row.TotalUserCount = d.TotalUserCount - } - } - out := make([]*pb.DBAppStat, 0, len(bucket)) - for _, r := range bucket { - out = append(out, r) - } - sort.Slice(out, func(i, j int) bool { return out[i].StatDate < out[j].StatDate }) - return out -} - func dbErr(err error, where string) *pb.ErrorData { return &pb.ErrorData{ Code: pb.ErrorCode_DBError, diff --git a/apps/services/modules/api/model.go b/apps/services/modules/api/model.go index 0570e237..3b95eb4e 100644 --- a/apps/services/modules/api/model.go +++ b/apps/services/modules/api/model.go @@ -11,7 +11,6 @@ import ( "yunyan/lego/sys/mysql" "yunyan/lego/sys/postgres" "yunyan/pb" - "yunyan/sys/appstat" "gorm.io/gorm" ) @@ -65,6 +64,10 @@ func (this *modelComp) Init(service core.IService, module core.IModule, comp cor if err = postgres.CreateTable(comm.TableSolutionProvider, &pb.DBSolutionProvider{}); err != nil { this.module.Errorln(err) + } else { + //方案商 id 从 0xAB01 起。用 Floor 而非 Start:表里已经有 1、2、3… 这些历史行时 + //Start 会直接跳过,新建的方案商仍然接着 4 发号;Floor 保留历史行、把序列抬到 0xAB01。 + postgres.AutoIncrementFloor(comm.TableSolutionProvider, "id", 0xAB01) } if err = postgres.CreateTable(comm.TableBrand, &pb.DBBrand{}); err != nil { this.module.Errorln(err) @@ -106,34 +109,9 @@ func (this *modelComp) Init(service core.IService, module core.IModule, comp cor //设置wakeupvoice表的主键从0xD001开始(仅空表生效,跨 MySQL/Postgres 通用) postgres.AutoIncrementStart(comm.TableWakeupVoice, "id", 0xD001) } - if err = mysql.CreateTable(comm.TableAppStatDaily, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatMonthly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatYearly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatGlobal, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } if err = mysql.CreateTable(comm.TableProductStat, &pb.DBProductStat{}); err != nil { this.module.Errorln(err) } - // AutoMigrate 不会主动放宽已存在列的 varchar 宽度, - // 历史库 stat_date 可能是更短的 varchar(N<10),导致写入 "YYYY-MM-DD" 报 1406 Data too long。 - // 这里强制把 4 张表的 stat_date 拉到 varchar(20),幂等。 - for _, tbl := range []string{ - comm.TableAppStatDaily, - comm.TableAppStatMonthly, - comm.TableAppStatYearly, - comm.TableAppStatGlobal, - } { - if e := mysql.Exec(fmt.Sprintf("ALTER TABLE %s MODIFY COLUMN stat_date VARCHAR(20) NOT NULL", tbl)).Error; e != nil { - this.module.Errorln(e) - } - } model := &pb.DBAdminUser{ Account: this.module.options.AdninAccount, Password: this.module.options.AdninPassword, @@ -906,25 +884,6 @@ func (this *modelComp) addIntegralLog(log *pb.DBUserUseLog) (err error) { return } -// 查询App综合统计(日/月/年) -func (this *modelComp) getAppStats(statType int32, startDate, endDate string) (stats []*pb.DBAppStat, err error) { - tbl := appstat.TableByType(statType) - stats = make([]*pb.DBAppStat, 0) - err = mysql.Find(tbl, &stats, "stat_date >= ? AND stat_date <= ?", startDate, endDate) - return -} - -// 查询全局累计统计(单行) -func (this *modelComp) getGlobalStat() (stat *pb.DBAppStat, err error) { - stat = &pb.DBAppStat{StatDate: "all"} - if err = mysql.FindOne(comm.TableAppStatGlobal, stat, "stat_date=?", "all"); err != nil { - // 不存在则返回空行,不视为错误 - stat = &pb.DBAppStat{StatDate: "all"} - err = nil - } - return -} - // 累计注册人数(用户表 COUNT(*)) func (this *modelComp) countTotalUsers() (total int64, err error) { err = mysql.Table(comm.TableUser).Count(&total).Error diff --git a/apps/services/modules/api/model_rebuildstats.go b/apps/services/modules/api/model_rebuildstats.go index bc5b1f55..8d27d025 100644 --- a/apps/services/modules/api/model_rebuildstats.go +++ b/apps/services/modules/api/model_rebuildstats.go @@ -8,211 +8,6 @@ import ( "time" ) -// 全量重建统计所需的源表聚合查询。 -// 注:以下分组使用 MySQL DATE(FROM_UNIXTIME(ts)) 按服务器时区切日, -// 与事件驱动写入 time.Now().Format("2006-01-02") 的本地时区一致。 - -// orderCreateAgg 创建订单数(按 create_time 切日) -type dateCountRow struct { - D string `gorm:"column:d"` - Count int64 `gorm:"column:cnt"` -} - -// dateSumPayRow 已支付订单 SUM(amount)/COUNT(*)/COUNT(DISTINCT uid) (按 pay_time 切日) -type dateSumPayRow struct { - D string `gorm:"column:d"` - Count int64 `gorm:"column:cnt"` - Amount int64 `gorm:"column:amt"` - UserCount int64 `gorm:"column:ucnt"` -} - -// dateUseLogRow 资源发放(按 ts 切日,4 个发放字段 SUM) -type dateUseLogRow struct { - D string `gorm:"column:d"` - AddVipDay int64 `gorm:"column:vipday"` - AddAgentInt int64 `gorm:"column:aichat"` - AddTradeSec int64 `gorm:"column:tradesec"` - AddMeetSec int64 `gorm:"column:meetsec"` -} - -// aggOrderCreate 按创建日聚合 order_count -func (this *modelComp) aggOrderCreate() (rows []*dateCountRow, err error) { - rows = make([]*dateCountRow, 0) - err = mysql.Table(comm.TablePayOrder). - Select("DATE_FORMAT(FROM_UNIXTIME(create_time), '%Y-%m-%d') AS d, COUNT(*) AS cnt"). - Where("create_time > 0"). - Group("d"). - Scan(&rows).Error - return -} - -// aggOrderFailed 按创建日聚合 failed_order_count(失败订单按 create_time 切日) -func (this *modelComp) aggOrderFailed() (rows []*dateCountRow, err error) { - rows = make([]*dateCountRow, 0) - err = mysql.Table(comm.TablePayOrder). - Select("DATE_FORMAT(FROM_UNIXTIME(create_time), '%Y-%m-%d') AS d, COUNT(*) AS cnt"). - Where("create_time > 0 AND status = ?", int32(pb.PayOrderStatus_PAY_ORDER_FAILED)). - Group("d"). - Scan(&rows).Error - return -} - -// aggOrderPaid 按支付日聚合 pay_amount/pay_count/pay_user_count -func (this *modelComp) aggOrderPaid() (rows []*dateSumPayRow, err error) { - rows = make([]*dateSumPayRow, 0) - err = mysql.Table(comm.TablePayOrder). - Select("DATE_FORMAT(FROM_UNIXTIME(pay_time), '%Y-%m-%d') AS d, COUNT(*) AS cnt, COALESCE(SUM(amount),0) AS amt, COUNT(DISTINCT uid) AS ucnt"). - Where("pay_time > 0 AND status = ?", int32(pb.PayOrderStatus_PAY_ORDER_PAID)). - Group("d"). - Scan(&rows).Error - return -} - -// aggUserSignup 按注册日聚合 new_user_count -func (this *modelComp) aggUserSignup() (rows []*dateCountRow, err error) { - rows = make([]*dateCountRow, 0) - err = mysql.Table(comm.TableUser). - Select("DATE_FORMAT(FROM_UNIXTIME(createtime), '%Y-%m-%d') AS d, COUNT(*) AS cnt"). - Where("createtime > 0"). - Group("d"). - Scan(&rows).Error - return -} - -// aggUseLogGrant 按日聚合资源发放(vipday/aichat/trade-sec/meet-sec) -func (this *modelComp) aggUseLogGrant() (rows []*dateUseLogRow, err error) { - rows = make([]*dateUseLogRow, 0) - err = mysql.Table(comm.TableUserUseLog). - Select("DATE_FORMAT(FROM_UNIXTIME(ts), '%Y-%m-%d') AS d, " + - "COALESCE(SUM(addvipday),0) AS vipday, " + - "COALESCE(SUM(addagentintegral),0) AS aichat, " + - "COALESCE(SUM(addtradesecond),0) AS tradesec, " + - "COALESCE(SUM(addmeetsecond),0) AS meetsec"). - Where("ts > 0"). - Group("d"). - Scan(&rows).Error - return -} - -// 全局累计字段(无法按日还原,只能 SUM 整库) - -// sumLifetimePayments 全部已支付订单累计金额、订单数、付费人数(去重 uid) -func (this *modelComp) sumLifetimePayments() (payAmount, payCount, payUserCount, orderCount, failedOrderCount int64, err error) { - type r struct { - Amount int64 `gorm:"column:amt"` - Count int64 `gorm:"column:cnt"` - UCnt int64 `gorm:"column:ucnt"` - } - var rr r - if err = mysql.Table(comm.TablePayOrder). - Select("COALESCE(SUM(amount),0) AS amt, COUNT(*) AS cnt, COUNT(DISTINCT uid) AS ucnt"). - Where("status = ?", int32(pb.PayOrderStatus_PAY_ORDER_PAID)). - Scan(&rr).Error; err != nil { - return - } - payAmount, payCount, payUserCount = rr.Amount, rr.Count, rr.UCnt - if err = mysql.Table(comm.TablePayOrder).Count(&orderCount).Error; err != nil { - return - } - err = mysql.Table(comm.TablePayOrder). - Where("status = ?", int32(pb.PayOrderStatus_PAY_ORDER_FAILED)). - Count(&failedOrderCount).Error - return -} - -// sumLifetimeUseLog 全部资源发放累计(用于 global 行 add_* 字段) -func (this *modelComp) sumLifetimeUseLog() (vipDay, aiInt, tradeSec, meetSec int64, err error) { - type r struct { - Vipday int64 `gorm:"column:vipday"` - Aichat int64 `gorm:"column:aichat"` - Tradesec int64 `gorm:"column:tradesec"` - Meetsec int64 `gorm:"column:meetsec"` - } - var rr r - err = mysql.Table(comm.TableUserUseLog). - Select("COALESCE(SUM(addvipday),0) AS vipday, " + - "COALESCE(SUM(addagentintegral),0) AS aichat, " + - "COALESCE(SUM(addtradesecond),0) AS tradesec, " + - "COALESCE(SUM(addmeetsecond),0) AS meetsec"). - Scan(&rr).Error - vipDay, aiInt, tradeSec, meetSec = rr.Vipday, rr.Aichat, rr.Tradesec, rr.Meetsec - return -} - -// sumLifetimeUserStats 从 userstatistics 表汇总累计消耗(用于 global 行的 ai/trade/meet 字段) -func (this *modelComp) sumLifetimeUserStats() (aiChat, aiUp, aiDown, tradeCount, tradeTime, tradeWords, meetCount, meetTime int64, err error) { - type r struct { - AiChat int64 `gorm:"column:aichat"` - AiUp int64 `gorm:"column:aiup"` - AiDown int64 `gorm:"column:aidown"` - TradeCount int64 `gorm:"column:tcnt"` - TradeTime int64 `gorm:"column:ttime"` - TradeWords int64 `gorm:"column:twords"` - MeetCount int64 `gorm:"column:mcnt"` - MeetTime int64 `gorm:"column:mtime"` - } - var rr r - err = mysql.Table(comm.TableUserStatistics). - Select("COALESCE(SUM(aichatnum),0) AS aichat, " + - "COALESCE(SUM(aichatuptoken),0) AS aiup, " + - "COALESCE(SUM(aichatdowntoken),0) AS aidown, " + - "COALESCE(SUM(" + colTradeNumSum + "),0) AS tcnt, " + - "COALESCE(SUM(" + colTradeTimeSum + "),0) AS ttime, " + - "COALESCE(SUM(tradewordcount),0) AS twords, " + - "COALESCE(SUM(meetnum),0) AS mcnt, " + - "COALESCE(SUM(meettime),0) AS mtime"). - Scan(&rr).Error - aiChat, aiUp, aiDown = rr.AiChat, rr.AiUp, rr.AiDown - tradeCount, tradeTime, tradeWords = rr.TradeCount, rr.TradeTime, rr.TradeWords - meetCount, meetTime = rr.MeetCount, rr.MeetTime - return -} - -// countUserDevices 累计绑定设备数 -func (this *modelComp) countUserDevices() (total int64, err error) { - err = mysql.Table(comm.TableUserdevice).Count(&total).Error - return -} - -// countActiveVipUsers 当前仍有效的 VIP 用户数(vipexptime > now);用于 global 行 vip_user_count 字段 -func (this *modelComp) countActiveVipUsers(nowTs int64) (total int64, err error) { - err = mysql.Table(comm.TableUser).Where("vipexptime > ?", nowTs).Count(&total).Error - return -} - -// truncateStatTables 清空 4 张统计表(在 stat_consumer 锁内调用) -func (this *modelComp) truncateStatTables() (err error) { - for _, tbl := range []string{ - comm.TableAppStatDaily, - comm.TableAppStatMonthly, - comm.TableAppStatYearly, - comm.TableAppStatGlobal, - } { - if err = mysql.Exec("TRUNCATE TABLE " + tbl).Error; err != nil { - return - } - } - return -} - -// insertStatRows 批量插入(已 TRUNCATE 后调用) -func (this *modelComp) insertStatRows(tbl string, rows []*pb.DBAppStat) (err error) { - if len(rows) == 0 { - return - } - const batch = 500 - for i := 0; i < len(rows); i += batch { - j := i + batch - if j > len(rows) { - j = len(rows) - } - if err = mysql.Table(tbl).CreateInBatches(rows[i:j], batch).Error; err != nil { - return - } - } - return -} - // rebuildUserStatisticsFromLogs 从 useruselog + echomeet_record 全量重建 userstatistics。 // // Phase 1 — useruselog(logtype=UserConsume): diff --git a/apps/services/modules/api/module.go b/apps/services/modules/api/module.go index bc21152a..51e7dcd6 100644 --- a/apps/services/modules/api/module.go +++ b/apps/services/modules/api/module.go @@ -16,7 +16,6 @@ type API struct { api *apiComp model *modelComp modelAudit *modelAuditComp - statConsumer *statConsumerComp permIntercept *permissionInterceptorComp options *Options } @@ -49,6 +48,5 @@ func (this *API) OnInstallComp() { this.api = this.RegisterComp(new(apiComp)).(*apiComp) this.model = this.RegisterComp(new(modelComp)).(*modelComp) this.modelAudit = this.RegisterComp(new(modelAuditComp)).(*modelAuditComp) - this.statConsumer = this.RegisterComp(new(statConsumerComp)).(*statConsumerComp) this.permIntercept = this.RegisterComp(new(permissionInterceptorComp)).(*permissionInterceptorComp) } diff --git a/apps/services/modules/api/stat_consumer.go b/apps/services/modules/api/stat_consumer.go deleted file mode 100644 index bb7e8df7..00000000 --- a/apps/services/modules/api/stat_consumer.go +++ /dev/null @@ -1,193 +0,0 @@ -package api - -import ( - "context" - "yunyan/comm" - "yunyan/lego/core" - "yunyan/lego/core/cbase" - "yunyan/lego/sys/log" - "yunyan/lego/sys/mysql" - redissys "yunyan/lego/sys/redis" - "yunyan/pb" - "yunyan/sys/appstat" - "encoding/json" - "sync" - "time" -) - -// statConsumerComp 消费 Redis 统计队列并写入日/月/年三张 DB 表。 -// 内存中缓存当前日/月/年三行数据,启动时从 DB 加载一次; -// 每次 flush 在内存累加后立即写入 DB,日期切换时加载新行。 -type statConsumerComp struct { - cbase.ModuleCompBase - module *API - cancel context.CancelFunc - - // 内存缓存(三个时间粒度 + 全局累计) - // mu 保护以下四个缓存指针的读写,避免 consumer flush 与外部全量重建并发改写。 - mu sync.Mutex - daily *pb.DBAppStat - monthly *pb.DBAppStat - yearly *pb.DBAppStat - global *pb.DBAppStat -} - -// 全局累计行固定 stat_date 标识 -const GlobalStatDate = "all" - -func (this *statConsumerComp) Init(service core.IService, module core.IModule, comp core.IModuleComp, opt core.IModuleOptions) (err error) { - this.ModuleCompBase.Init(service, module, comp, opt) - this.module = module.(*API) - return -} - -func (this *statConsumerComp) Start() (err error) { - if err = this.ModuleCompBase.Start(); err != nil { - return - } - // 启动时从 DB 加载当前日/月/年缓存行 - this.loadCache(time.Now()) - - ctx, cancel := context.WithCancel(context.Background()) - this.cancel = cancel - go this.consume(ctx) - log.Infoln("appstat consumer started") - return -} - -func (this *statConsumerComp) Destroy() (err error) { - if this.cancel != nil { - this.cancel() - } - return this.ModuleCompBase.Destroy() -} - -// loadCache 从 DB 加载指定时间点对应的日/月/年/全局行(不存在则初始化空行) -func (this *statConsumerComp) loadCache(t time.Time) { - this.daily = this.loadOrNew(comm.TableAppStatDaily, t.Format("2006-01-02")) - this.monthly = this.loadOrNew(comm.TableAppStatMonthly, t.Format("2006-01")) - this.yearly = this.loadOrNew(comm.TableAppStatYearly, t.Format("2006")) - this.global = this.loadOrNew(comm.TableAppStatGlobal, GlobalStatDate) -} - -func (this *statConsumerComp) loadOrNew(tbl, date string) *pb.DBAppStat { - var row pb.DBAppStat - if err := mysql.FindOne(tbl, &row, "stat_date=?", date); err != nil { - row = pb.DBAppStat{StatDate: date} - } - return &row -} - -// consume 从 Redis 队列消费统计增量 -func (this *statConsumerComp) consume(ctx context.Context) { - for { - select { - case <-ctx.Done(): - return - default: - } - vals, err := redissys.Conn().BRPop(ctx, 5*time.Second, redissys.RKey(appstat.RedisQueueKey)).Result() - if err != nil { - continue - } - if len(vals) < 2 { - continue - } - var delta pb.DBAppStat - if err = json.Unmarshal([]byte(vals[1]), &delta); err != nil { - this.module.Warn("stat consumer unmarshal", log.Field{Key: "err", Value: err.Error()}) - continue - } - this.flush(&delta) - } -} - -// checkExpired 各粒度独立检查是否过期,过期则创建新的空统计对象;全局行不过期 -func (this *statConsumerComp) checkExpired(now time.Time) { - if today := now.Format("2006-01-02"); this.daily == nil || this.daily.StatDate != today { - this.daily = &pb.DBAppStat{StatDate: today} - } - if month := now.Format("2006-01"); this.monthly == nil || this.monthly.StatDate != month { - this.monthly = &pb.DBAppStat{StatDate: month} - } - if year := now.Format("2006"); this.yearly == nil || this.yearly.StatDate != year { - this.yearly = &pb.DBAppStat{StatDate: year} - } - if this.global == nil { - this.global = this.loadOrNew(comm.TableAppStatGlobal, GlobalStatDate) - } -} - -// flush 将 delta 累加到内存缓存后立即写入 DB -func (this *statConsumerComp) flush(delta *pb.DBAppStat) { - if len(delta.StatDate) < 10 { - this.module.Warn("stat consumer flush: invalid StatDate", log.Field{Key: "date", Value: delta.StatDate}) - return - } - this.mu.Lock() - defer this.mu.Unlock() - this.checkExpired(time.Now()) - - accum(this.daily, delta) - accum(this.monthly, delta) - accum(this.yearly, delta) - accum(this.global, delta) - - this.saveRow(comm.TableAppStatDaily, this.daily) - this.saveRow(comm.TableAppStatMonthly, this.monthly) - this.saveRow(comm.TableAppStatYearly, this.yearly) - this.saveRow(comm.TableAppStatGlobal, this.global) -} - -// ReloadCache 由外部(如全量重建)在改写 DB 后调用,丢弃过期内存缓存重新从 DB 加载。 -func (this *statConsumerComp) ReloadCache() { - this.mu.Lock() - defer this.mu.Unlock() - this.loadCache(time.Now()) -} - -// WithLock 提供给外部对四张统计表执行原子改写(清空+重建)使用, -// 期间 consumer flush 被阻塞,避免被半成品状态覆盖。 -func (this *statConsumerComp) WithLock(fn func() error) error { - this.mu.Lock() - defer this.mu.Unlock() - return fn() -} - -func (this *statConsumerComp) saveRow(tbl string, row *pb.DBAppStat) { - if row == nil { - return - } - row.UpdateTime = time.Now().Unix() - if err := mysql.Save(tbl, row); err != nil { - this.module.Warn("stat consumer save", log.Field{Key: "table", Value: tbl}, log.Field{Key: "err", Value: err.Error()}) - } -} - -// accum 将 delta 各字段累加到 row 上 -func accum(row, delta *pb.DBAppStat) { - row.PayAmount += delta.PayAmount - row.OrderCount += delta.OrderCount - row.PayCount += delta.PayCount - row.FailedOrderCount += delta.FailedOrderCount - row.PayUserCount += delta.PayUserCount - row.AddVipDays += delta.AddVipDays - row.AddAiIntegral += delta.AddAiIntegral - row.AddTradeSecond += delta.AddTradeSecond - row.AddMeetSecond += delta.AddMeetSecond - row.AiChatCount += delta.AiChatCount - row.AiUpToken += delta.AiUpToken - row.AiDownToken += delta.AiDownToken - row.TradeCount += delta.TradeCount - row.TradeTime += delta.TradeTime - row.TradeWords += delta.TradeWords - row.MeetCount += delta.MeetCount - row.MeetTime += delta.MeetTime - row.NewUserCount += delta.NewUserCount - row.LoginCount += delta.LoginCount - row.ActiveUserCount += delta.ActiveUserCount - row.TotalUserCount = delta.TotalUserCount // 快照值,直接覆盖 - row.VipUserCount += delta.VipUserCount - row.BindDeviceCount += delta.BindDeviceCount - row.ActiveDeviceCount += delta.ActiveDeviceCount -} diff --git a/apps/services/modules/console/registry.go b/apps/services/modules/console/registry.go index fa8cf7aa..ea417b73 100644 --- a/apps/services/modules/console/registry.go +++ b/apps/services/modules/console/registry.go @@ -109,6 +109,9 @@ func ensureDeviceTables() error { if err := adb.CreateTable(comm.TableSolutionProvider, &pb.DBSolutionProvider{}); err != nil { return err } + // 方案商 id 从 0xAB01 起:Floor 语义,已有历史行(id 1、2、3…)保留,只把序列抬上去, + // 后台新建的方案商从 0xAB01 开始发号。 + _ = adb.AutoIncrementFloor(comm.TableSolutionProvider, "id", 0xAB01) if err := adb.CreateTable(comm.TableBrand, &pb.DBBrand{}); err != nil { return err } diff --git a/apps/services/modules/pay/model.go b/apps/services/modules/pay/model.go index 9f331289..83b19d50 100644 --- a/apps/services/modules/pay/model.go +++ b/apps/services/modules/pay/model.go @@ -21,15 +21,6 @@ func (this *modelComp) Init(service core.IService, module core.IModule, comp cor if err = mysql.CreateTable(comm.TablePayOrder, &pb.DBPayOrder{}); err != nil { this.module.Errorln(err) } - if err = mysql.CreateTable(comm.TableAppStatDaily, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatMonthly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatYearly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } return } diff --git a/apps/services/modules/user/api_v2_getagents.go b/apps/services/modules/user/api_v2_getagents.go index 3da82ff0..d4bea8a7 100644 --- a/apps/services/modules/user/api_v2_getagents.go +++ b/apps/services/modules/user/api_v2_getagents.go @@ -1,7 +1,6 @@ package user import ( - "os" "sort" "yunyan/comm" @@ -22,7 +21,7 @@ import ( // 注册在 apiV2Comp(Version=v2)上,路由 = user_getagents_v2;已在 EncryptMsgs 中声明, // 网关对整个响应体做 AES-CBC 加密(系统提示词/编排关系不明文出网关)。 func (this *apiV2Comp) GetAgents(session comm.IUserSession, req *pb.UserGetAgentsReq) (resp *pb.UserGetAgentsResp, errdata *pb.ErrorData) { - app := os.Getenv("APP_NAME") + app := comm.AppName() rows := make([]*comm.AgentConfig, 0) if err := postgres.Find(comm.TableAgentConfig, &rows, "(app_name=? OR app_name='')", app); err != nil && err != postgres.ErrNoDocuments { diff --git a/apps/services/modules/user/api_v2_getthirdsvcs.go b/apps/services/modules/user/api_v2_getthirdsvcs.go index 9157026b..6156f894 100644 --- a/apps/services/modules/user/api_v2_getthirdsvcs.go +++ b/apps/services/modules/user/api_v2_getthirdsvcs.go @@ -25,7 +25,7 @@ import ( // 网关对整个响应体做 AES-CBC 加密——响应含第三方服务明文凭据,绝不能明文出网关。 func (this *apiV2Comp) GetThirdSvcs(session comm.IUserSession, req *pb.UserGetThirdSvcsReq) (resp *pb.UserGetThirdSvcsResp, errdata *pb.ErrorData) { var ( - app = os.Getenv("APP_NAME") + app = comm.AppName() encKey = os.Getenv("FIELD_ENCRYPT_KEY") region = pb.Region_RegionUSA ) diff --git a/apps/services/modules/user/model_user.go b/apps/services/modules/user/model_user.go index 8b324739..00aeaac5 100644 --- a/apps/services/modules/user/model_user.go +++ b/apps/services/modules/user/model_user.go @@ -10,7 +10,6 @@ import ( redissys "yunyan/lego/sys/redis" "yunyan/pb" "fmt" - "os" "time" "gorm.io/gorm" @@ -39,15 +38,6 @@ func (this *modelUserComp) Init(service core.IService, module core.IModule, comp if err = mysql.CreateTable(comm.TableUserUseLog, &pb.DBUserUseLog{}); err != nil { this.module.Errorln(err) } - if err = mysql.CreateTable(comm.TableAppStatDaily, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatMonthly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } - if err = mysql.CreateTable(comm.TableAppStatYearly, &pb.DBAppStat{}); err != nil { - this.module.Errorln(err) - } return } @@ -237,8 +227,10 @@ func (this *modelUserComp) getProductVersions(pid uint32) (models []*pb.DBProduc // 下发同语义。下发给客户端的仍是 pb.DBChannelApp,由 toPbChannelApp 做映射。 // loadChannelApps 读本应用作用域 + 全局默认的全部渠道配置行。 +// 应用名走 comm.AppName()(APP_NAME 优先、回退 ANALYZE_APP_NAME),否则只填了统计应用名的 +// 部署会读不到自己那一层配置,只剩全局默认行。 func (this *modelUserComp) loadChannelApps() ([]*comm.ChannelApp, error) { - app := os.Getenv("APP_NAME") + app := comm.AppName() rows := make([]*comm.ChannelApp, 0) err := postgres.Find(comm.TableChannelAppCfg, &rows, "(app_name=? OR app_name='')", app) if err != nil && err != postgres.ErrNoDocuments { diff --git a/apps/services/sys/appstat/appstat.go b/apps/services/sys/appstat/appstat.go deleted file mode 100644 index 74fa8034..00000000 --- a/apps/services/sys/appstat/appstat.go +++ /dev/null @@ -1,47 +0,0 @@ -// Package appstat 提供 App 综合统计的生产者接口。 -// -// 各模块调用 appstat.Incr(delta) 将统计增量序列化后推入 Redis 队列, -// 由 api 模块的 statConsumerComp 负责消费并写入 DB。 -package appstat - -import ( - "context" - "yunyan/comm" - "yunyan/lego/sys/log" - redissys "yunyan/lego/sys/redis" - "yunyan/pb" - "encoding/json" - "time" -) - -const RedisQueueKey = "appstat:queue" - -// TableByType 根据统计类型返回对应表名(1=日 2=月 3=年) -func TableByType(statType int32) string { - switch statType { - case 2: - return comm.TableAppStatMonthly - case 3: - return comm.TableAppStatYearly - default: - return comm.TableAppStatDaily - } -} - -// Incr 将统计增量序列化后推入 Redis 队列,非阻塞。 -// delta.StatDate 必须为 "2006-01-02" 格式。 -func Incr(delta *pb.DBAppStat) { - if delta.UpdateTime == 0 { - delta.UpdateTime = time.Now().Unix() - } - data, err := json.Marshal(delta) - if err != nil { - log.Warn("appstat.Incr marshal", log.Field{Key: "err", Value: err.Error()}) - return - } - ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) - defer cancel() - if err = redissys.Conn().LPush(ctx, redissys.RKey(RedisQueueKey), data).Err(); err != nil { - log.Warn("appstat.Incr lpush", log.Field{Key: "err", Value: err.Error()}) - } -} diff --git a/deploy/app/env/env.example b/deploy/app/env/env.example index 78d398b5..f24d470d 100644 --- a/deploy/app/env/env.example +++ b/deploy/app/env/env.example @@ -35,6 +35,9 @@ FIELD_ENCRYPT_KEY= # api 的应用服务配置托管、echomeet 的会议记录服务编排都按 (APP_NAME, APP_REGION) 作用域解析; # 不配则落到「全局默认」作用域('',0)——编排的全局默认层仍然生效,但按应用/区域的细分配置不生效。 # APP_REGION 用区域短代码:cn/us/ea/sea/...(见 comm/region.go)。 +# 注:按应用作用域读配置的接口(渠道分发 user_getchannelapp(s)、user_getagents_v2、 +# user_getthirdsvcs_v2)走 comm.AppName()——APP_NAME 留空时回退下面的 ANALYZE_APP_NAME, +# 两者填一个即可;但 api 的「服务配置托管」只认 APP_NAME,要用那个功能就必须填这里。 APP_NAME= APP_REGION=