You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

1249 lines
39 KiB

package console
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"path/filepath"
"strconv"
"strings"
"sync"
"time"
"yunyan/comm"
"yunyan/lego/core"
"yunyan/lego/core/cbase"
"yunyan/lego/sys/log"
"yunyan/lego/sys/mysql"
"yunyan/pb"
"github.com/gin-gonic/gin"
"github.com/golang-jwt/jwt/v4"
)
// serverComp 控制面对外 HTTP 服务:静态页 + 注册表 API + /web/api/api_*(按 X-App-Id 直连应用库)+ 配置事件下发。
type serverComp struct {
cbase.ModuleCompBase
module *Console
options *Options
srv *http.Server
// analyzing 记录正在跑历史分析的 app_id(值无意义),防止同一应用并发分析——
// analyzeAppHistory 含 DELETE+重写,并发会互相清表,必须串行。重复点击直接拒绝。
analyzing sync.Map
// deploying 记录正在升级/部署的部署 id(值无意义),防止同一部署并发 SSH 部署——
// docker compose up 并发会互相拉起/删容器,必须串行。重复点击直接拒绝。
deploying sync.Map
}
func (this *serverComp) 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.(*Console)
this.options = opt.(*Options)
gin.SetMode(gin.ReleaseMode)
eng := gin.New()
eng.Use(gin.Recovery())
eng.Use(cors())
this.routes(eng)
this.srv = &http.Server{Addr: this.options.HTTP.Addr, Handler: eng}
return
}
func (this *serverComp) Start() (err error) {
if err = this.ModuleCompBase.Start(); err != nil {
return
}
go func() {
if e := this.srv.ListenAndServe(); e != nil && e != http.ErrServerClosed {
log.Errorf("console: HTTP 服务退出: %v", e)
}
}()
log.Infof("console: HTTP 服务监听 %s(静态目录 %s)", this.options.HTTP.Addr, this.options.StaticDir)
return nil
}
func (this *serverComp) Destroy() (err error) {
if this.srv != nil {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = this.srv.Shutdown(ctx)
}
return this.ModuleCompBase.Destroy()
}
func (this *serverComp) routes(eng *gin.Engine) {
eng.GET("/console/health", func(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"code": 0, "msg": "ok"})
})
eng.GET("/", func(c *gin.Context) { c.Redirect(http.StatusFound, "/console/") })
// 兼容现有前端调用约定:/web/api/api_<method>。login/getsiteinfo 放行,其余需登录态。
web := eng.Group("/web/api", this.authWeb())
web.POST("/:method", this.handleWeb)
// /console/api:需登录态,再按角色细分。
mgr := eng.Group("/console/api", this.authStrict())
// 当前登录账号信息(侧边栏底部展示:名称/身份/代理资源点余额)——任意登录角色可取自己的。
mgr.POST("/me", this.accountMe)
// 应用注册表 CRUD + 配置事件:超管 + 管理员(代理商/运营不能改应用接入)。
apps := mgr.Group("", this.requireIdentity(pb.Identity_Admin, pb.Identity_Manager))
// 应用/产品(按应用名分组的父层):超管 + 管理员。
apps.POST("/products/list", this.productsList)
apps.POST("/products/add", this.productsAdd)
apps.POST("/products/update", this.productsUpdate)
apps.POST("/products/del", this.productsDel)
apps.POST("/apps/list", this.appsList)
apps.POST("/apps/test", this.appsTest)
apps.POST("/apps/analyze", this.appsAnalyze)
apps.POST("/apps/upgrade", this.appsUpgrade) // 远程升级:SSH 到部署环境 → docker login/pull/up 重启(NDJSON 流式回显)
apps.POST("/apps/tags", this.appsTags) // 列镜像仓库可用 tag(升级弹窗下拉选版本)
apps.POST("/apps/add", this.appsAdd)
apps.POST("/apps/update", this.appsUpdate)
apps.POST("/apps/del", this.appsDel)
apps.POST("/apps/assign", this.appsAssign)
apps.POST("/apps/bindenv", this.appsBindEnv)
// 应用×服务 基础设施/运行配置:应用服务配置页读写(默认值由业务服务启动 seed)。
apps.POST("/svcconf/list", this.svcRuntimeList)
apps.POST("/svcconf/services", this.svcRuntimeServices)
apps.POST("/svcconf/save", this.svcRuntimeSave)
apps.POST("/svcconf/del", this.svcRuntimeDel)
// 应用业务功能配置(登录/短信/邮件/支付…):读写应用业务库(app_module_config),非公共库。
apps.POST("/modulecfg/catalog", this.moduleCfgCatalog)
apps.POST("/modulecfg/list", this.moduleCfgList)
apps.POST("/modulecfg/save", this.moduleCfgSave)
apps.POST("/modulecfg/del", this.moduleCfgDel)
// 部署环境(区域)管理:与应用注册同权限(超管 + 管理员)。
apps.POST("/envs/list", this.envList)
apps.POST("/envs/test", this.envTest)
apps.POST("/envs/add", this.envAdd)
apps.POST("/envs/update", this.envUpdate)
apps.POST("/envs/del", this.envDel)
apps.POST("/config/notify", this.configNotify)
// 给代理账号充值/调整资源点余额:超管 + 管理员(与赠送给代理同义)。
apps.POST("/accounts/topup", this.accountsTopup)
// 账号管理:仅超管(pb.Identity 注明 Admin 为唯一可管理后台账号者)。
acc := mgr.Group("/accounts", this.requireIdentity(pb.Identity_Admin))
acc.POST("/list", this.accountsList)
acc.POST("/add", this.accountsAdd)
acc.POST("/update", this.accountsUpdate)
acc.POST("/del", this.accountsDel)
acc.POST("/resetpwd", this.accountsResetPwd)
// 静态后台页用 NoRoute 兜底:API 路由优先,未命中再当 /console/ 下的静态文件,
// 避免 gin 的 catch-all 通配与 /console/api、/console/health 冲突。
eng.NoRoute(this.serveStatic)
}
// serveStatic 服务 StaticDir 下 /console/ 前缀的静态资源(含 SPA 入口 index.html)。
func (this *serverComp) serveStatic(c *gin.Context) {
p := c.Request.URL.Path
if !strings.HasPrefix(p, "/console/") {
c.Status(http.StatusNotFound)
return
}
rel := strings.TrimPrefix(p, "/console/")
if rel == "" || strings.HasSuffix(rel, "/") {
rel += "index.html"
}
if strings.Contains(rel, "..") { // 防目录穿越
c.Status(http.StatusBadRequest)
return
}
c.File(filepath.Join(this.options.StaticDir, filepath.Clean(rel)))
}
// ============================ 鉴权中间件 ============================
// authWeb 用于 /web/api:放行 api_login / api_getsiteinfo,其余校验 JWT。
func (this *serverComp) authWeb() gin.HandlerFunc {
return func(c *gin.Context) {
method := c.Param("method")
if method == "api_login" || method == "api_getsiteinfo" {
c.Next()
return
}
if !this.tokenValid(c) {
c.AbortWithStatusJSON(http.StatusOK, &comm.HttpResult{Code: pb.ErrorCode_NoLogin, Message: "未登录"})
return
}
c.Next()
}
}
// authStrict 用于注册表/事件接口:一律校验 JWT。
func (this *serverComp) authStrict() gin.HandlerFunc {
return func(c *gin.Context) {
if !this.tokenValid(c) {
c.AbortWithStatusJSON(http.StatusOK, &comm.HttpResult{Code: pb.ErrorCode_NoLogin, Message: "未登录"})
return
}
c.Next()
}
}
func (this *serverComp) tokenValid(c *gin.Context) bool {
tok := c.GetHeader("Authorization")
if tok == "" {
return false
}
claims, err := parseToken(tok, []byte(this.options.TokenKey))
if err != nil {
return false
}
// 解析出的身份与作用域存进上下文,供 requireIdentity / scopeOf 与后续 handler 取用。
c.Set("identity", claims.Identity)
c.Set("username", claims.Username)
c.Set("account_id", claims.AccountId)
c.Set("apps", claims.Apps)
c.Set("products", claims.Products)
c.Set("regions", claims.Regions)
return true
}
// currentIdentity 取当前登录者角色(需在 authWeb/authStrict 之后);未登录返回 Null。
func currentIdentity(c *gin.Context) pb.Identity {
if v, ok := c.Get("identity"); ok {
if idt, ok := v.(pb.Identity); ok {
return idt
}
}
return pb.Identity_Identity_Null
}
// currentAccountId 取当前登录者的后台账号 id(yaml 引导超管为 0)。
func currentAccountId(c *gin.Context) uint32 {
if v, ok := c.Get("account_id"); ok {
if id, ok := v.(uint32); ok {
return id
}
}
return 0
}
// requireIdentity 角色守卫:当前登录者角色须在 allowed 集合内,否则 403(无权限)。
// 用在 authStrict 之后(依赖其写入的 identity)。
func (this *serverComp) requireIdentity(allowed ...pb.Identity) gin.HandlerFunc {
return func(c *gin.Context) {
idt := currentIdentity(c)
for _, a := range allowed {
if idt == a {
c.Next()
return
}
}
c.AbortWithStatusJSON(http.StatusOK, &comm.HttpResult{Code: pb.ErrorCode_InsufficientPermissions, Message: "无权限"})
}
}
// ============================ /web/api 派发 ============================
func (this *serverComp) handleWeb(c *gin.Context) {
method := c.Param("method")
switch method {
case "api_login":
this.login(c)
return
case "api_getsiteinfo":
this.getSiteInfo(c)
return
// 统计 dashboard:数据在 console 自己的公共库,不需选中应用,先于设备域分派。
case "api_getstatssummary":
this.getStatsSummary(c)
return
case "api_getstatstrend":
this.getStatsTrend(c)
return
// 统计日志:console 从 NATS 收到的埋点接收流水,数据在 console 自己的 Redis,与选中应用无关。
case "api_getstatlogs":
this.getStatLogs(c)
return
case "api_clearstatlogs":
this.clearStatLogs(c)
return
// COS 上传凭证:全局配置,与选中应用无关。
case "api_getcostoken":
this.getCosToken(c)
return
// 服务配置域:全局环境配置/MCP/会议模板,三类共享配置存 console 主库(supabase),与选中应用无关。
case "api_getglobalconfigs":
this.getGlobalConfigs(c)
return
case "api_addglobalconfig":
this.addGlobalConfig(c)
return
case "api_updateglobalconfig":
this.updateGlobalConfig(c)
return
case "api_delglobalconfig":
this.delGlobalConfig(c)
return
case "api_getmcpservers":
this.getMcpServers(c)
return
case "api_addmcpserver":
this.addMcpServer(c)
return
case "api_updatemcpserver":
this.updateMcpServer(c)
return
case "api_delmcpserver":
this.delMcpServer(c)
return
case "api_getmeettemplates":
this.getMeetTemplates(c)
return
case "api_getmeettemplate":
this.getMeetTemplate(c)
return
case "api_addmeettemplate":
this.addMeetTemplate(c)
return
case "api_updatemeettemplate":
this.updateMeetTemplate(c)
return
case "api_delmeettemplate":
this.delMeetTemplate(c)
return
// 第三方服务配置(svcconfig):服务字段定义/默认值 + 区域字段分叉(值覆盖/专属字段/停用字段) + 多区域同步,存 console 主库,与选中应用无关。
case "api_getsvcconfigs":
this.getSvcConfigs(c)
return
case "api_addsvcconfig":
this.addSvcConfig(c)
return
case "api_updatesvcconfig":
this.updateSvcConfig(c)
return
case "api_delsvcconfig":
this.delSvcConfig(c)
return
case "api_getsvcregionoverride":
this.getSvcRegionOverride(c)
return
case "api_savesvcregionoverride":
this.saveSvcRegionOverride(c)
return
case "api_syncsvctoregins":
this.syncSvcToRegions(c)
return
// 第三方服务商模板(svctemplate):数据驱动的字段 schema 预设,新增服务时选模板自动带出字段。
case "api_getsvctemplates":
this.getSvcTemplates(c)
return
case "api_addsvctemplate":
this.addSvcTemplate(c)
return
case "api_updatesvctemplate":
this.updateSvcTemplate(c)
return
case "api_delsvctemplate":
this.delSvcTemplate(c)
return
// yunyan 自管 Agent 配置(agent):编排第三方服务(引用 svc id),配置下发客户端;存 console 主库,与选中应用无关。
case "api_getagents":
this.getAgents(c)
return
case "api_getagent":
this.getAgent(c)
return
case "api_addagent":
this.addAgent(c)
return
case "api_updateagent":
this.updateAgent(c)
return
case "api_delagent":
this.delAgent(c)
return
// 通话翻译配置(calltranslate):区域×语言对×服务编排,存 console 主库,与选中应用无关。
case "api_getcalltranslatepairs":
this.getCallTranslatePairs(c)
return
case "api_savecalltranslatepair":
this.saveCallTranslatePair(c)
return
case "api_delcalltranslatepair":
this.delCallTranslatePair(c)
return
case "api_synccalltranslatepairs":
this.syncCallTranslatePairs(c)
return
// 邮件下发(生产设备码清单):multipart 表单 + 全局 email 系统,与选中应用无关。
case "api_sendsmtpemail":
this.sendSmtpEmail(c)
return
// 用户查询:按选中应用(X-App-Id)直连其业务库查 user 表,与设备域(supabase)无关。
case "api_getuserinfo":
this.getUserInfo(c)
return
// 赠送资源点:给选中应用业务库的终端用户加 VIP天/翻译分/会议分;代理扣自己余额。
case "api_giftuserresource":
this.giftUserResource(c)
return
}
// 设备管理域:厂家/产品/设备码/公码/出货单都在 console 主库(supabase),无需选应用。
sys := consoleDeviceConn()
// 仅重置类要清各应用业务库的 userdevice/product_stat,按 X-App-Id 取该应用业务库填入 sys.service。
switch method {
case "api_resetlicensestatus", "api_batchresetlicensestatus", "api_restfactorydevics":
appId := parseAppId(c.GetHeader("X-App-Id"))
if !this.requireAppScope(c, appId) {
return
}
svc, err := this.module.registry.getServiceDB(appId)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
sys.service = svc
}
switch method {
case "api_getfactorys":
this.getFactorys(c, sys)
case "api_getproducts":
this.getProducts(c, sys)
case "api_getproduct":
this.getProduct(c, sys)
case "api_getproductversions":
this.getProductVersions(c, sys)
case "api_addproduct":
this.addProduct(c, sys)
case "api_updateproduct":
this.updateProduct(c, sys)
case "api_delproduct":
this.delProduct(c, sys)
case "api_addproductversion":
this.addProductVersion(c, sys)
case "api_delproductversion":
this.delProductVersion(c, sys)
case "api_addfactory":
this.addFactory(c, sys)
case "api_updatefactory":
this.updateFactory(c, sys)
case "api_delfactory":
this.delFactory(c, sys)
case "api_createfactorydevics":
this.createFactoryDevics(c, sys)
case "api_disablefactorydevics":
this.disableFactoryDevics(c, sys)
case "api_revokefactorydevics":
this.revokeFactoryDevics(c, sys)
case "api_getfactorydevics":
this.getFactoryDevics(c, sys)
case "api_getfactorydevic":
this.getFactoryDevic(c, sys)
case "api_getfactorydeliverynotes":
this.getFactoryDeliveryNotes(c, sys)
case "api_resetlicensestatus":
this.resetLicenseStatus(c, sys)
case "api_batchresetlicensestatus":
this.batchResetLicenseStatus(c, sys)
case "api_restfactorydevics":
this.restFactoryDevics(c, sys)
case "api_getfactorypubliccodes":
this.getFactoryPublicCodes(c, sys)
case "api_createfactorypubliccode":
this.createFactoryPublicCode(c, sys)
case "api_updatefactorypubliccode":
this.updateFactoryPublicCode(c, sys)
case "api_delfactorypubliccode":
this.delFactoryPublicCode(c, sys)
default:
writeErr(c, pb.ErrorCode_NoFindServiceHandleFunc, "无此接口: "+method)
}
}
func parseAppId(s string) uint32 {
n, _ := strconv.ParseUint(strings.TrimSpace(s), 10, 32)
return uint32(n)
}
// requireAppScope 校验 X-App-Id 指向的应用在当前账号作用域内(代理/运营)。
// 越界即写好错误响应并返回 false,调用方直接 return。超管/管理员或未绑定应用维度时直接放行。
func (this *serverComp) requireAppScope(c *gin.Context, appId uint32) bool {
s := scopeOf(c)
if s.unlimited || len(s.apps) == 0 {
return true
}
app, err := this.module.model.getApp(appId)
if err != nil || app == nil {
writeErr(c, pb.ErrorCode_ReqParameterError, "应用不存在或无权访问")
return false
}
if !s.allowApp(app.Name) {
writeErr(c, pb.ErrorCode_InsufficientPermissions, "无权访问该应用")
return false
}
return true
}
// ============================ 登录 / 站点信息 ============================
func (this *serverComp) login(c *gin.Context) {
var req pb.ApiLoginReq
_ = c.ShouldBindJSON(&req)
identity, accountId, ok := this.authenticate(req.Account, req.Password)
if !ok {
writeErr(c, pb.ErrorCode_ReqParameterError, "账号或密码错误")
return
}
now := time.Now()
// 作用域绑定(代理/运营才有;超管/管理员为空=不受限):先查出来,既写进 token 供请求级强制隔离,也回传前端控菜单/筛选。
var access, apps, regions, products string
if accountId > 0 {
if acc, e := this.module.model.getAccount(accountId); e == nil && acc != nil {
access, apps, regions, products = acc.Access, acc.Apps, acc.Regions, acc.Products
}
}
claims := &consoleClaims{
Identity: identity,
Username: req.Account,
AccountId: accountId,
Apps: apps,
Products: products,
Regions: regions,
RegisteredClaims: jwt.RegisteredClaims{
Issuer: "console",
Subject: fmt.Sprintf("%d", identity),
ExpiresAt: jwt.NewNumericDate(now.Add(24 * time.Hour)),
NotBefore: jwt.NewNumericDate(now),
IssuedAt: jwt.NewNumericDate(now),
ID: req.Account,
},
}
ts, err := jwt.NewWithClaims(jwt.SigningMethodHS256, claims).SignedString([]byte(this.options.TokenKey))
if err != nil {
writeErr(c, pb.ErrorCode_SystemError, err.Error())
return
}
// 把绑定应用名解析为 {id,name}:代理无权调 apps/list,但 users.html / 切换器需要注册 id 取 X-App-Id。
appBinds := make([]gin.H, 0)
if apps != "" {
if all, e := this.module.model.listApps(); e == nil {
want := map[string]bool{}
for _, n := range strings.Split(apps, ",") {
if n = strings.TrimSpace(n); n != "" {
want[n] = true
}
}
for _, a := range all {
if want[a.Name] {
appBinds = append(appBinds, gin.H{"id": a.Id, "name": a.Name, "region": a.Region})
}
}
}
}
writeOK(c, gin.H{
"account": req.Account,
"identity": identity,
"token": ts,
"avatar": "https://gw.alipayobjects.com/zos/rmsportal/BiazfanxmamNRoxxVxka.png",
"access": access,
"apps": apps,
"regions": regions,
"products": products,
"appbinds": appBinds, // [{id,name,region}]:代理用它取应用注册 id
})
}
// authenticate 校验账号密码,返回角色与账号 id。
// 顺序:先匹配 yaml 引导超管(AdminAccount/AdminPassword,account_id=0,保证首登可用);
// 再查 console_account 表(bcrypt 校验,须 Enabled)。
func (this *serverComp) authenticate(username, password string) (pb.Identity, uint32, bool) {
if username == this.options.AdminAccount && password == this.options.AdminPassword {
return pb.Identity_Admin, 0, true
}
acc, err := this.module.model.getAccountByName(username)
if err != nil || acc == nil || acc.Id == 0 || !acc.Enabled {
return pb.Identity_Identity_Null, 0, false
}
if !checkPassword(acc.Password, password) {
return pb.Identity_Identity_Null, 0, false
}
return acc.Identity, acc.Id, true
}
func (this *serverComp) getSiteInfo(c *gin.Context) {
writeOK(c, &pb.ApiGetSiteInfoResp{Sitename: this.options.SiteName, CurrencyType: 0})
}
// ============================ 设备管理域 handlers ============================
//
// 逻辑从 modules/api 的对应 handler 移植;去掉 RPCX session 与 scope 白名单(第一期仅超管),
// 数据访问改为作用于传入的"选中应用"连接 sys。
// 读 handler 统一走 deviceCache(缓存优先 + DB 回退),设备域引用数据恒命中 Redis,避免直连海外主库的卡顿。
func (this *serverComp) getFactorys(c *gin.Context, sys *appConn) {
models, err := this.module.deviceCache.GetFactorys()
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
// 代理/运营:由绑定产品反推可见厂家(只保留拥有绑定产品的厂家)。
if s := scopeOf(c); !s.unlimited && len(s.products) > 0 {
allowFid := make(map[uint32]bool)
if prods, perr := this.module.deviceCache.GetProducts(); perr == nil {
for _, p := range prods {
if s.allowProduct(p.Id) {
allowFid[p.Factoryid] = true
}
}
}
kept := make([]*pb.DBFactory, 0, len(models))
for _, f := range models {
if allowFid[f.Id] {
kept = append(kept, f)
}
}
models = kept
}
writeOK(c, &pb.ApiGetFactorysResp{Factorys: models})
}
func (this *serverComp) getProducts(c *gin.Context, sys *appConn) {
models, err := this.module.deviceCache.GetProducts()
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
models = scopeOf(c).filterProducts(models) // 代理/运营:仅返回绑定产品
writeOK(c, &pb.ApiGetProductsResp{Products: models})
}
func (this *serverComp) getProduct(c *gin.Context, sys *appConn) {
var req pb.ApiGetProductReq
_ = c.ShouldBindJSON(&req)
if !scopeOf(c).allowProduct(req.Id) {
writeErr(c, pb.ErrorCode_InsufficientPermissions, "无权访问该产品")
return
}
model, err := this.module.deviceCache.GetProduct(req.Id)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiGetProductResp{Product: model})
}
func (this *serverComp) getProductVersions(c *gin.Context, sys *appConn) {
var req pb.ApiGetProductVersionsReq
_ = c.ShouldBindJSON(&req)
if !scopeOf(c).allowProduct(req.Pid) {
writeErr(c, pb.ErrorCode_InsufficientPermissions, "无权访问该产品")
return
}
models, err := this.module.deviceCache.GetProductVersions(req.Pid)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiGetProductVersionsResp{Versions: models})
}
func (this *serverComp) delFactory(c *gin.Context, sys *appConn) {
var req pb.ApiDelFactoryReq
_ = c.ShouldBindJSON(&req)
if err := dvDelFactory(sys, req.Id); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
this.module.deviceCache.refreshFactorys() // 写后刷新厂家缓存
writeOK(c, &pb.ApiDelFactoryResp{})
}
func (this *serverComp) getFactoryDevics(c *gin.Context, sys *appConn) {
var req pb.ApiGetFactoryDevicsReq
_ = c.ShouldBindJSON(&req)
if req.Productid == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "productid 不能为空")
return
}
var (
conds = make([]string, 0, 7)
args []interface{}
)
if req.Code != "" {
conds = append(conds, "code=?")
args = append(args, req.Code)
}
if req.Devicemac != "" {
conds = append(conds, "devicemac=?")
args = append(args, req.Devicemac)
}
if req.Factoryid != 0 {
conds = append(conds, "factoryid=?")
args = append(args, req.Factoryid)
}
if req.Probatch != 0 {
conds = append(conds, "probatch=?")
args = append(args, req.Probatch)
}
// status: -1 仅未使用, 1 已使用, 2 已过期, 0 不筛选
switch req.Status {
case -1:
conds = append(conds, "status=?")
args = append(args, 0)
case 1, 2:
conds = append(conds, "status=?")
args = append(args, req.Status)
}
if req.Uid != "" {
conds = append(conds, "uid=?")
args = append(args, req.Uid)
}
query := ""
if len(conds) > 0 {
query = strings.Join(conds, " and ")
}
sortBy := strings.ToLower(strings.TrimSpace(req.Sortby))
if sortBy != "usedtime" && sortBy != "createtime" {
sortBy = "createtime"
}
order := strings.ToLower(strings.TrimSpace(req.Order))
if order != "asc" {
order = "desc"
}
limit := int(req.Limit)
if limit <= 0 {
limit = 200
}
if limit > 2000 {
limit = 2000
}
offset := int(req.Offset)
if offset < 0 {
offset = 0
}
orderClause := fmt.Sprintf("%s %s", sortBy, order)
models, total, err := dvFactoryDevicsPaged(sys, req.Productid, query, orderClause, limit, offset, args...)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiGetFactoryDevicsResp{Devices: models, Total: total})
}
func (this *serverComp) getFactoryDevic(c *gin.Context, sys *appConn) {
var req pb.ApiGetFactoryDevicReq
_ = c.ShouldBindJSON(&req)
model, err := dvFactoryDevic(sys, req.Productid, req.Devicemac)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiGetFactoryDevicResp{Device: model})
}
func (this *serverComp) getFactoryDeliveryNotes(c *gin.Context, sys *appConn) {
var req pb.ApiGetFactoryDeliveryNotesReq
_ = c.ShouldBindJSON(&req)
query := "1=1"
args := make([]interface{}, 0)
if req.Factoryid > 0 {
query += " AND factoryid = ?"
args = append(args, req.Factoryid)
}
if req.Productid > 0 {
query += " AND productid = ?"
args = append(args, req.Productid)
}
if req.Start > 0 {
query += " AND ts >= ?"
args = append(args, req.Start)
}
if req.End > 0 {
query += " AND ts <= ?"
args = append(args, req.End)
}
limit := int(req.Limit)
if limit <= 0 || limit > 2000 {
limit = 500
}
query += " ORDER BY ts DESC LIMIT ?"
args = append(args, limit)
notes, err := dvFactoryDeliveryNotes(sys, query, args...)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
var totalDevices int64
for _, n := range notes {
totalDevices += int64(n.Devicetnum)
}
writeOK(c, &pb.ApiGetFactoryDeliveryNotesResp{
Notes: notes,
TotalBatch: int64(len(notes)),
TotalDevices: totalDevices,
})
}
func (this *serverComp) resetLicenseStatus(c *gin.Context, sys *appConn) {
var req pb.ApiResetLicenseStatusReq
_ = c.ShouldBindJSON(&req)
// 走完整重置(删 userdevice + 清 license + 扣激活统计),与批量复用同一 resetOneLicense。
if err := resetOneLicense(sys, req.License); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiResetLicenseStatusResp{})
}
func (this *serverComp) getFactoryPublicCodes(c *gin.Context, sys *appConn) {
var req pb.ApiGetFactoryPublicCodesReq
_ = c.ShouldBindJSON(&req)
models, err := dvFactoryPublicCodes(sys, req.Factoryid, req.Status)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiGetFactoryPublicCodesResp{Codes: models})
}
func (this *serverComp) createFactoryPublicCode(c *gin.Context, sys *appConn) {
var req pb.ApiCreateFactoryPublicCodeReq
_ = c.ShouldBindJSON(&req)
if req.Factoryid == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "factoryid 必填")
return
}
if _, err := dvFactory(sys, req.Factoryid); err != nil {
writeErr(c, pb.ErrorCode_DBError, "厂家不存在: "+err.Error())
return
}
// 防碰撞:最多重试 5 次
var code string
for i := 0; i < 5; i++ {
gc, err := comm.GeneratePublicCode()
if err != nil {
writeErr(c, pb.ErrorCode_SystemError, err.Error())
return
}
if _, e := dvFactoryPublicCode(sys, gc); e != nil { // 不存在 → 可用
code = gc
break
}
}
if code == "" {
writeErr(c, pb.ErrorCode_SystemError, "公码生成多次碰撞,请重试")
return
}
now := time.Now().Unix()
model := &pb.DBFactoryPublicCode{
Code: code,
Factoryid: req.Factoryid,
Productid: req.Productid,
Remark: req.Remark,
Status: 0,
Maxuses: req.Maxuses,
Useduses: 0,
Expiretime: req.Expiretime,
Bindrewardvpitime: req.Bindrewardvpitime,
Bindrewardtranslate: req.Bindrewardtranslate,
Bindrewardmeeting: req.Bindrewardmeeting,
Createtime: now,
Updatetime: now,
}
if err := dvAddFactoryPublicCode(sys, model); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiCreateFactoryPublicCodeResp{Code: model})
}
func (this *serverComp) updateFactoryPublicCode(c *gin.Context, sys *appConn) {
var req pb.ApiUpdateFactoryPublicCodeReq
_ = c.ShouldBindJSON(&req)
if req.Code == "" {
writeErr(c, pb.ErrorCode_ReqParameterError, "code 必填")
return
}
model, err := dvFactoryPublicCode(sys, req.Code)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, "公码不存在: "+err.Error())
return
}
model.Productid = req.Productid
model.Remark = req.Remark
model.Status = req.Status
model.Maxuses = req.Maxuses
model.Expiretime = req.Expiretime
model.Bindrewardvpitime = req.Bindrewardvpitime
model.Bindrewardtranslate = req.Bindrewardtranslate
model.Bindrewardmeeting = req.Bindrewardmeeting
model.Updatetime = time.Now().Unix()
if err = dvSaveFactoryPublicCode(sys, model); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiUpdateFactoryPublicCodeResp{Code: model})
}
func (this *serverComp) delFactoryPublicCode(c *gin.Context, sys *appConn) {
var req pb.ApiDelFactoryPublicCodeReq
_ = c.ShouldBindJSON(&req)
if req.Code == "" {
writeErr(c, pb.ErrorCode_ReqParameterError, "code 必填")
return
}
if err := dvDelFactoryPublicCode(sys, req.Code); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &pb.ApiDelFactoryPublicCodeResp{})
}
// ============================ 用户查询 ============================
// getUserInfo 用户查询:前端选定「应用 + 区域」即定位到唯一注册项(X-App-Id),
// 据此经 registry.getServiceDB 取该部署的业务库(MySQL)连接,按 uid 直查 user 表返回。
// 业务数据在各应用自己的业务库,故不走 console 主库(supabase)。
func (this *serverComp) getUserInfo(c *gin.Context) {
appId := parseAppId(c.GetHeader("X-App-Id"))
if !this.requireAppScope(c, appId) {
return
}
conn, err := this.module.registry.getServiceDB(appId)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
var req struct {
Uid string `json:"uid"`
}
_ = c.ShouldBindJSON(&req)
uid := strings.TrimSpace(req.Uid)
if uid == "" {
writeErr(c, pb.ErrorCode_ReqParameterError, "请输入要查询的用户UID")
return
}
user := &pb.DBUser{}
if err := conn.FindOne(comm.TableUser, user, "uid=?", uid); err != nil {
if errors.Is(err, mysql.ErrNoDocuments) {
writeErr(c, pb.ErrorCode_UserSessionNobeing, "未找到该用户")
return
}
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
user.Password = "" // 不下发密码
writeOK(c, user)
}
// giftUserResource 赠送资源点给终端用户:直连选中应用(X-App-Id)业务库,给 user 加
// VIP天/翻译分钟/会议分钟,并写 useruselog(AdminGive) 流水。
// - 超管 / 管理员:无消耗,直接发放;
// - 代理:先原子扣减自己的资源点余额,不足即拒绝;写库失败回滚已扣余额;
// - 运营:无权赠送。
func (this *serverComp) giftUserResource(c *gin.Context) {
identity := currentIdentity(c)
if identity != pb.Identity_Admin && identity != pb.Identity_Manager && identity != pb.Identity_Agent {
writeErr(c, pb.ErrorCode_InsufficientPermissions, "当前角色无权赠送资源点")
return
}
var req struct {
Uid string `json:"uid"`
Vipday int64 `json:"vipday"`
Trademin int64 `json:"trademin"`
Meetmin int64 `json:"meetmin"`
}
_ = c.ShouldBindJSON(&req)
uid := strings.TrimSpace(req.Uid)
if uid == "" {
writeErr(c, pb.ErrorCode_ReqParameterError, "请输入要赠送的用户UID")
return
}
if req.Vipday < 0 || req.Trademin < 0 || req.Meetmin < 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "赠送数量不能为负")
return
}
if req.Vipday == 0 && req.Trademin == 0 && req.Meetmin == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "请至少填写一项赠送数量")
return
}
appId := parseAppId(c.GetHeader("X-App-Id"))
if !this.requireAppScope(c, appId) {
return
}
conn, err := this.module.registry.getServiceDB(appId)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
user := &pb.DBUser{}
if err := conn.FindOne(comm.TableUser, user, "uid=?", uid); err != nil {
if errors.Is(err, mysql.ErrNoDocuments) {
writeErr(c, pb.ErrorCode_UserSessionNobeing, "未找到该用户")
return
}
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
now := time.Now().Unix()
accId := currentAccountId(c)
// 代理:先原子扣减自己的资源点余额(CAS 防超发),不足即拒绝。
if identity == pb.Identity_Agent {
ok, e := this.module.model.agentDeductBalance(accId, req.Vipday, req.Trademin, req.Meetmin, now)
if e != nil {
writeErr(c, pb.ErrorCode_DBError, e.Error())
return
}
if !ok {
writeErr(c, pb.ErrorCode_InsufficientPermissions, "您的资源点余额不足")
return
}
}
// 应用到用户余额(VIP 过期时间:已过期/未开通则从现在起算,否则顺延;与公码奖励口径一致)。
addTradeSec := req.Trademin * 60
addMeetSec := req.Meetmin * 60
if req.Vipday > 0 {
addVipSec := req.Vipday * 24 * 60 * 60
if user.Vipexptime == 0 || user.Vipexptime < now {
user.Vipexptime = now + addVipSec
} else {
user.Vipexptime += addVipSec
}
}
if addTradeSec > 0 {
user.Tradeintegral += addTradeSec
user.Tradetotalintegral += addTradeSec
}
if addMeetSec > 0 {
user.Meetintegral += addMeetSec
user.Meettotalintegral += addMeetSec
}
if err := conn.Save(comm.TableUser, user); err != nil {
// 代理:写库失败回滚已扣余额,避免用户没到账却扣了代理。
if identity == pb.Identity_Agent {
_, _ = this.module.model.agentAddBalance(accId, req.Vipday, req.Trademin, req.Meetmin, now)
}
writeErr(c, pb.ErrorCode_DBError, "赠送失败: "+err.Error())
return
}
// 流水(业务库 useruselog)。已到账,失败不回滚,仅吞掉。
_ = conn.Insert(comm.TableUserUseLog, &pb.DBUserUseLog{
Uid: uid,
Ts: now,
Logtype: pb.UserLogType_AdminGive,
Addvipday: req.Vipday,
Addtradesecond: addTradeSec,
Addmeetsecond: addMeetSec,
Extra: fmt.Sprintf("console gift by identity=%d account=%d", int32(identity), accId),
})
resp := gin.H{
"uid": uid,
"vipexptime": user.Vipexptime,
"tradeintegral": user.Tradeintegral,
"meetintegral": user.Meetintegral,
"tradetotalintegral": user.Tradetotalintegral,
"meettotalintegral": user.Meettotalintegral,
}
// 代理:回传扣减后的最新余额,前端据此刷新侧边栏。
if identity == pb.Identity_Agent {
if acc, e := this.module.model.getAccount(accId); e == nil && acc != nil {
resp["balance_vipday"] = acc.BalanceVipday
resp["balance_trademin"] = acc.BalanceTrademin
resp["balance_meetmin"] = acc.BalanceMeetmin
}
}
writeOK(c, resp)
}
// accountMe 返回当前登录账号的展示信息(侧边栏底部用):用户名、角色、代理资源点余额。
func (this *serverComp) accountMe(c *gin.Context) {
username, _ := c.Get("username")
resp := gin.H{"username": username, "identity": currentIdentity(c)}
if accId := currentAccountId(c); accId > 0 {
if acc, err := this.module.model.getAccount(accId); err == nil && acc != nil {
resp["balance_vipday"] = acc.BalanceVipday
resp["balance_trademin"] = acc.BalanceTrademin
resp["balance_meetmin"] = acc.BalanceMeetmin
}
}
writeOK(c, resp)
}
// accountsTopup 管理员给代理账号充值 / 调整资源点余额(按类增量,可正可负;任一类不能调成负数)。
func (this *serverComp) accountsTopup(c *gin.Context) {
var req struct {
AccountId uint32 `json:"account_id"`
Vipday int64 `json:"vipday"`
Trademin int64 `json:"trademin"`
Meetmin int64 `json:"meetmin"`
}
if err := c.ShouldBindJSON(&req); err != nil {
writeErr(c, pb.ErrorCode_ReqParameterError, err.Error())
return
}
if req.AccountId == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "缺少账号id")
return
}
if req.Vipday == 0 && req.Trademin == 0 && req.Meetmin == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "请至少填写一项调整数量")
return
}
acc, err := this.module.model.getAccount(req.AccountId)
if err != nil || acc == nil || acc.Id == 0 {
writeErr(c, pb.ErrorCode_DBError, "账号不存在")
return
}
if acc.Identity != pb.Identity_Agent {
writeErr(c, pb.ErrorCode_ReqParameterError, "只能给代理账号充值资源点")
return
}
ok, e := this.module.model.agentAddBalance(req.AccountId, req.Vipday, req.Trademin, req.Meetmin, time.Now().Unix())
if e != nil {
writeErr(c, pb.ErrorCode_DBError, e.Error())
return
}
if !ok {
writeErr(c, pb.ErrorCode_ReqParameterError, "调整后余额不能为负")
return
}
acc, _ = this.module.model.getAccount(req.AccountId)
writeOK(c, gin.H{
"balance_vipday": acc.BalanceVipday,
"balance_trademin": acc.BalanceTrademin,
"balance_meetmin": acc.BalanceMeetmin,
})
}
// ============================ 注册表 CRUD ============================
func (this *serverComp) appsList(c *gin.Context) {
models, err := this.module.model.listApps()
if err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, models)
}
func (this *serverComp) appsAdd(c *gin.Context) {
var model AppRegistry
if err := c.ShouldBindJSON(&model); err != nil {
writeErr(c, pb.ErrorCode_ReqParameterError, err.Error())
return
}
if msg := checkDeployment(&model); msg != "" {
writeErr(c, pb.ErrorCode_ReqParameterError, msg)
return
}
model.Id = 0
now := time.Now().Unix()
model.Createtime = now
model.Updatetime = now
if err := this.module.model.addApp(&model); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, &model)
}
func (this *serverComp) appsUpdate(c *gin.Context) {
var model AppRegistry
if err := c.ShouldBindJSON(&model); err != nil {
writeErr(c, pb.ErrorCode_ReqParameterError, err.Error())
return
}
if model.Id == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填")
return
}
if msg := checkDeployment(&model); msg != "" {
writeErr(c, pb.ErrorCode_ReqParameterError, msg)
return
}
model.Updatetime = time.Now().Unix()
if err := this.module.model.saveApp(&model); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
this.module.registry.invalidate(model.Id) // 配置变了,失效连接缓存
writeOK(c, &model)
}
func (this *serverComp) appsDel(c *gin.Context) {
var req struct {
Id uint32 `json:"id"`
}
if err := c.ShouldBindJSON(&req); err != nil || req.Id == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填")
return
}
if err := this.module.model.delApp(req.Id); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
this.module.registry.invalidate(req.Id)
writeOK(c, gin.H{"id": req.Id})
}
// appsAssign 轻量「归类」:只把某部署行的 app_name 改成目标应用(空=移回未分组),不做连通性复检。
// 用于把历史部署行归到应用名下,避免走完整表单+测试连接。
func (this *serverComp) appsAssign(c *gin.Context) {
var req struct {
Id uint32 `json:"id"`
AppName string `json:"app_name"`
}
if err := c.ShouldBindJSON(&req); err != nil || req.Id == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "id 必填")
return
}
model, err := this.module.model.getApp(req.Id)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, "部署不存在: "+err.Error())
return
}
model.AppName = req.AppName
model.Updatetime = time.Now().Unix()
if err := this.module.model.saveApp(model); err != nil {
writeErr(c, pb.ErrorCode_DBError, err.Error())
return
}
writeOK(c, model)
}
// ============================ 配置事件下发 ============================
func (this *serverComp) configNotify(c *gin.Context) {
var req struct {
AppId uint32 `json:"app_id"`
Type string `json:"type"`
Action string `json:"action"`
Payload json.RawMessage `json:"payload"`
}
if err := c.ShouldBindJSON(&req); err != nil || req.AppId == 0 {
writeErr(c, pb.ErrorCode_ReqParameterError, "app_id 必填")
return
}
app, err := this.module.model.getApp(req.AppId)
if err != nil {
writeErr(c, pb.ErrorCode_DBError, "应用未注册: "+err.Error())
return
}
ev := &ConfigEvent{
AppId: req.AppId,
Type: req.Type,
Action: req.Action,
Payload: req.Payload,
Ts: time.Now().Unix(),
}
if err := publishConfigEvent(app, ev); err != nil {
writeErr(c, pb.ErrorCode_SystemError, "事件发布失败: "+err.Error())
return
}
writeOK(c, gin.H{"subject": eventSubject(app)})
}