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.
246 lines
10 KiB
246 lines
10 KiB
package console
|
|
|
|
import (
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
|
|
"yunyan/comm"
|
|
"yunyan/lego/core"
|
|
"yunyan/lego/core/cbase"
|
|
"yunyan/lego/sys/log"
|
|
"yunyan/lego/sys/mysql"
|
|
"yunyan/lego/sys/postgres"
|
|
"yunyan/pb"
|
|
)
|
|
|
|
// registryComp 应用业务库连接管理器:按 appId 把注册项的 ServiceDsn 物化成业务库连接并缓存。
|
|
//
|
|
// 设备库内容(厂家/产品/设备码/公码/出货单)统一在 console 主库(supabase),不分应用、无需选;
|
|
// 只有「重置设备码」要清各应用自己的 userdevice/product_stat 时,才按 X-App-Id 取应用业务库。
|
|
type registryComp struct {
|
|
cbase.ModuleCompBase
|
|
module *Console
|
|
mu sync.RWMutex
|
|
conns map[uint32]mysql.ISys // appId → 业务库(ServiceDsn) 连接
|
|
}
|
|
|
|
func (this *registryComp) 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.conns = make(map[uint32]mysql.ISys)
|
|
return
|
|
}
|
|
|
|
// getServiceDB 取(或惰性建)指定应用的业务库连接(用于清 userdevice/product_stat)。
|
|
func (this *registryComp) getServiceDB(appId uint32) (mysql.ISys, error) {
|
|
if appId == 0 {
|
|
return nil, fmt.Errorf("未指定应用(appId=0),重置设备码需先在切换器选择应用(清理该应用的用户绑定)")
|
|
}
|
|
this.mu.RLock()
|
|
if conn, ok := this.conns[appId]; ok {
|
|
this.mu.RUnlock()
|
|
return conn, nil
|
|
}
|
|
this.mu.RUnlock()
|
|
|
|
app, err := this.module.model.getApp(appId)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("应用未注册或读取失败(appId=%d): %w", appId, err)
|
|
}
|
|
if !app.Enabled {
|
|
return nil, fmt.Errorf("应用已停用(appId=%d, name=%s)", appId, app.Name)
|
|
}
|
|
|
|
this.mu.Lock()
|
|
defer this.mu.Unlock()
|
|
// 双重检查:可能在抢锁期间别的请求已建好。
|
|
if conn, ok := this.conns[appId]; ok {
|
|
return conn, nil
|
|
}
|
|
conn, err := openDriver(app.ServiceDsn)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("应用业务库连接失败(appId=%d, name=%s): %w", appId, app.Name, err)
|
|
}
|
|
// 建连后补齐 console 会直接读写的业务库新列:console 用最新的 pb 结构操作 user/userdevice/
|
|
// payorder,而这些表的迁移平时由【业务服务启动时 AutoMigrate】完成。只升级 console、
|
|
// 业务服务还是老镜像时,缺列会让赠送资源等接口直接报 Unknown column。这里兜一层。
|
|
ensureServiceColumns(conn, app.Name)
|
|
this.conns[appId] = conn
|
|
return conn, nil
|
|
}
|
|
|
|
// serviceColumn 一个需要在业务库补齐的列(表名 + 用于取 schema 的模型 + 列名)。
|
|
type serviceColumn struct {
|
|
table string
|
|
model any
|
|
col string
|
|
}
|
|
|
|
// serviceColumns console 直接读写业务库时依赖的、由渠道分成体系新增的列。
|
|
// 与 pb 结构一致;补列走 gorm Migrator,跨 MySQL/Postgres 通用(MySQL 无 ADD COLUMN IF NOT EXISTS)。
|
|
var serviceColumns = []serviceColumn{
|
|
{comm.TableUser, &pb.DBUser{}, "lastbindchannelid"},
|
|
{comm.TableUserdevice, &pb.DBUserDivice{}, "channelid"},
|
|
{comm.TablePayOrder, &pb.DBPayOrder{}, "src_license"},
|
|
{comm.TablePayOrder, &pb.DBPayOrder{}, "src_mac"},
|
|
{comm.TablePayOrder, &pb.DBPayOrder{}, "src_productid"},
|
|
{comm.TablePayOrder, &pb.DBPayOrder{}, "src_brandid"},
|
|
{comm.TablePayOrder, &pb.DBPayOrder{}, "src_channelid"},
|
|
}
|
|
|
|
// ensureServiceColumns 幂等补齐业务库缺失的列。全程 best-effort:
|
|
// 表不存在(该应用没这块业务)或补列失败都只告警,不阻断连接建立。
|
|
func ensureServiceColumns(conn mysql.ISys, appName string) {
|
|
for _, sc := range serviceColumns {
|
|
m := conn.Table(sc.table).Migrator()
|
|
if !m.HasTable(sc.table) {
|
|
continue
|
|
}
|
|
if m.HasColumn(sc.model, sc.col) {
|
|
continue
|
|
}
|
|
if err := m.AddColumn(sc.model, sc.col); err != nil {
|
|
log.Warn("console: 业务库补列失败",
|
|
log.Field{Key: "app", Value: appName},
|
|
log.Field{Key: "table", Value: sc.table},
|
|
log.Field{Key: "column", Value: sc.col},
|
|
log.Field{Key: "err", Value: err.Error()})
|
|
continue
|
|
}
|
|
log.Infof("console: 业务库 %s 补列 %s.%s 完成", appName, sc.table, sc.col)
|
|
}
|
|
}
|
|
|
|
// invalidate 失效某应用的业务库连接缓存(注册项被改/删时调用),下次访问按新配置重建。
|
|
func (this *registryComp) invalidate(appId uint32) {
|
|
this.mu.Lock()
|
|
delete(this.conns, appId)
|
|
this.mu.Unlock()
|
|
}
|
|
|
|
// ============================ 设备域连接 ============================
|
|
|
|
// appConn 设备管理域用的连接句柄:
|
|
// - device 恒为 console 主库(supabase)——厂家/产品/设备码/公码/出货单都在这里;
|
|
// - service 仅「重置设备码」时按选中应用填入其业务库——清 userdevice/product_stat。
|
|
//
|
|
// 暴露 AdminDB()/ServiceDB() 使 model_device.go 的 handler 调用方式保持不变。
|
|
type appConn struct {
|
|
device mysql.ISys // 设备库 = console 主库(supabase)
|
|
service mysql.ISys // 业务库 = 选中应用的 MySQL(仅重置类用到)
|
|
}
|
|
|
|
func (a *appConn) AdminDB() mysql.ISys { return a.device }
|
|
func (a *appConn) ServiceDB() mysql.ISys { return a.service }
|
|
|
|
// consoleDeviceConn 设备域默认连接:device 指向 console 主库(supabase),service 留空。
|
|
// 重置类操作再把 service 补成选中应用的业务库。
|
|
func consoleDeviceConn() *appConn {
|
|
return &appConn{device: postgres.GetSys()}
|
|
}
|
|
|
|
// openDriver 按 DSN 前缀选驱动:postgres:// 走 PostgreSQL,否则 MySQL(应用业务库通常为 MySQL)。
|
|
// 直接用子系统的 NewSys 工厂建连;两驱动 ISys 方法集一致,统一用 mysql.ISys 承载。
|
|
func openDriver(dsn string) (mysql.ISys, error) {
|
|
if strings.HasPrefix(dsn, "postgres://") || strings.HasPrefix(dsn, "postgresql://") {
|
|
return postgres.NewSys(postgres.SetDsn(dsn))
|
|
}
|
|
return mysql.NewSys(mysql.SetMySQLDsn(dsn))
|
|
}
|
|
|
|
// ensureDeviceTables 在 console 主库(supabase)幂等建好设备管理域所需的表(含按 product 的 license 分表)。
|
|
// 启动时跑一次:设备库内容已集中到 console 自己的库,不再按应用建。
|
|
func ensureDeviceTables() error {
|
|
adb := postgres.GetSys()
|
|
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
|
|
}
|
|
_ = adb.AutoIncrementStart(comm.TableBrand, "id", 0xA001)
|
|
// createtime/sharerate 是后加字段(品牌商建档时间、分成比例万分比):CreateTable 对已存在的表
|
|
// 跳过 AutoMigrate,故显式补列。历史品牌商 createtime=0(列表展示为空)、sharerate=0(不分成)。
|
|
for _, ddl := range []string{
|
|
"ALTER TABLE " + comm.TableBrand + " ADD COLUMN IF NOT EXISTS createtime bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableBrand + " ADD COLUMN IF NOT EXISTS sharerate integer DEFAULT 0",
|
|
} {
|
|
if res := adb.Exec(ddl); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
if err := adb.CreateTable(comm.TableChannel, &pb.DBChannel{}); err != nil {
|
|
return err
|
|
}
|
|
// official/sharerate 是后加字段(品牌官方渠道标记、渠道分成比例):CreateTable 对已存在的表
|
|
// 跳过 AutoMigrate,故显式补列。历史行 official=false(存量品牌商没有官方渠道,需要的话在
|
|
// 后台手工建一条)、sharerate=0(不分成)。
|
|
for _, ddl := range []string{
|
|
"ALTER TABLE " + comm.TableChannel + " ADD COLUMN IF NOT EXISTS official boolean DEFAULT false",
|
|
"ALTER TABLE " + comm.TableChannel + " ADD COLUMN IF NOT EXISTS sharerate integer DEFAULT 0",
|
|
} {
|
|
if res := adb.Exec(ddl); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
if err := adb.CreateTable(comm.TableProductionBatch, &pb.DBProductionBatch{}); err != nil {
|
|
return err
|
|
}
|
|
if err := adb.CreateTable(comm.TableProduct, &pb.DBProduct{}); err != nil {
|
|
return err
|
|
}
|
|
_ = adb.AutoIncrementStart(comm.TableProduct, "id", 0xB001)
|
|
// appnames 是后加字段(产品归属的应用名 CSV):CreateTable 对已存在的表跳过 AutoMigrate,故显式补列。
|
|
// 历史行留空,语义为「未绑定应用」——此时品牌商账号的应用下拉退化为不限制(见 brandAppNames)。
|
|
if res := adb.Exec("ALTER TABLE " + comm.TableProduct + " ADD COLUMN IF NOT EXISTS appnames varchar(255) DEFAULT ''"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
// providerid/createtime/description 同为后加字段(产品级绑定方案商、建档时间、产品说明),同样要显式补列。
|
|
// providerid 历史行为 0 = 未绑定方案商:生产设备码前必须在后台补选,否则 CID 取不到。
|
|
for _, ddl := range []string{
|
|
"ALTER TABLE " + comm.TableProduct + " ADD COLUMN IF NOT EXISTS providerid bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableProduct + " ADD COLUMN IF NOT EXISTS createtime bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableProduct + " ADD COLUMN IF NOT EXISTS description text DEFAULT ''",
|
|
// pb 里 providerid 带 gorm index;建表走 AutoMigrate 会自动建,补列这条路要自己补上
|
|
"CREATE INDEX IF NOT EXISTS idx_product_providerid ON " + comm.TableProduct + " (providerid)",
|
|
} {
|
|
if res := adb.Exec(ddl); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
if err := adb.CreateTable(comm.TableProductVersion, &pb.DBProductVersion{}); err != nil {
|
|
return err
|
|
}
|
|
if err := adb.CreateTable(comm.TableLicense, &pb.DBAuthCode{}); err != nil {
|
|
return err
|
|
}
|
|
if err := adb.CreateTable(comm.TableFactoryPublicCode, &pb.DBFactoryPublicCode{}); err != nil {
|
|
return err
|
|
}
|
|
// 按 product 列表建 license_%x 分表。
|
|
var products []*pb.DBProduct
|
|
if err := adb.Find(comm.TableProduct, &products, ""); err != nil && err != postgres.ErrNoDocuments {
|
|
return err
|
|
}
|
|
for _, p := range products {
|
|
if err := adb.CreateTable(fmt.Sprintf("%s_%x", comm.TableLicense, p.Id), &pb.DBAuthCode{}); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
// 渠道商维度是后加字段:CreateTable 对已存在的表跳过 AutoMigrate,故批次台账与全部 license 分表
|
|
// 都要显式补列(幂等)。历史行 channelid 留空,语义为「未分配渠道」。
|
|
tables := []string{comm.TableProductionBatch, comm.TableLicense}
|
|
for _, p := range products {
|
|
tables = append(tables, fmt.Sprintf("%s_%x", comm.TableLicense, p.Id))
|
|
}
|
|
for _, t := range tables {
|
|
if res := adb.Exec("ALTER TABLE " + t + " ADD COLUMN IF NOT EXISTS channelid varchar(16) DEFAULT ''"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|