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_。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 // 文件直传:签发阿里云 OSS 预签名 PUT URL(配置读自「第三方服务配置→存储→阿里云 OSS」)。 // 路由名沿用 api_getcostoken(前端 useApi.uploadFile 复用),实现见 api_upload.go。 case "api_getcostoken": this.getUploadURL(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)}) }