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.
402 lines
18 KiB
402 lines
18 KiB
package console
|
|
|
|
import (
|
|
"fmt"
|
|
"sort"
|
|
"strconv"
|
|
"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()
|
|
// 芯片厂商:方案商的 vendor 列指向它,故先于 solution_provider 建好。
|
|
if err := adb.CreateTable(comm.TableChipVendor, &ChipVendor{}); err != nil {
|
|
return err
|
|
}
|
|
// 补种原 SolutionVendor 枚举的 4 个取值(id 与枚举数值一一对应,存量方案商靠它归组)。
|
|
// 幂等:已存在的行不覆盖,运营改过的名字不会被启动时刷回去。
|
|
// 一次性迁移:上一版把 CID 做成独立列(id 自增 + cid 另存),本版改成「id 即 CID」。
|
|
// 已部署过上一版的库要把带 cid 的行搬到以 CID 为主键,并级联改方案商的 vendor 引用。
|
|
// 迁移完会删掉 cid 列,之后本调用直接跳过。必须先于补种执行。
|
|
if err := migrateChipVendorCidToId(); err != nil {
|
|
return err
|
|
}
|
|
if err := ensureChipVendorSeed(); err != nil {
|
|
return err
|
|
}
|
|
// id 即 CID、由人工指定,不再自增,故不设自增下限。
|
|
// 代工厂/组装厂档案(与品牌商并列的独立实体)。
|
|
// ⚠️ 表名是 oem_factory,不是库里那张无人引用的遗留 factory 表,别混。
|
|
if err := adb.CreateTable(comm.TableOemFactory, &OemFactory{}); err != nil {
|
|
return err
|
|
}
|
|
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。
|
|
// ⚠️ sharerate 已停止在后台维护,分成比例改按生产批次配(production_batch.sharerate);
|
|
// 列不删是因为月结算暂时还读它,且删列会丢掉存量比例。
|
|
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
|
|
}
|
|
// factoryid / sharerate 是后加字段(这批货由哪家代工厂生产、给渠道商的分成比例万分比):
|
|
// CreateTable 对已存在的表跳过 AutoMigrate,故显式补列。
|
|
// 存量批次 factoryid=0 = 未指定(建这批时还没有工厂档案);sharerate=0 = 不分成
|
|
// (这些批次是在「比例挂在渠道商上」的年代生产的,比例得回 channel.sharerate 看)。
|
|
for _, ddl := range []string{
|
|
"ALTER TABLE " + comm.TableProductionBatch + " ADD COLUMN IF NOT EXISTS factoryid bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableProductionBatch + " ADD COLUMN IF NOT EXISTS sharerate integer DEFAULT 0",
|
|
// 绑定赠送三项:存量批次全 0,绑定时会回退到产品上的旧字段,故不需要 backfill。
|
|
"ALTER TABLE " + comm.TableProductionBatch + " ADD COLUMN IF NOT EXISTS rewardvip bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableProductionBatch + " ADD COLUMN IF NOT EXISTS rewardtranslate bigint DEFAULT 0",
|
|
"ALTER TABLE " + comm.TableProductionBatch + " ADD COLUMN IF NOT EXISTS rewardmeeting bigint DEFAULT 0",
|
|
} {
|
|
if res := adb.Exec(ddl); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
if err := adb.CreateTable(comm.TableProduct, &pb.DBProduct{}); err != nil {
|
|
return err
|
|
}
|
|
_ = adb.AutoIncrementStart(comm.TableProduct, "id", 0xB001)
|
|
// appnames 是后加字段(产品归属的应用名 CSV):CreateTable 对已存在的表跳过 AutoMigrate,故显式补列。
|
|
// 历史行留空,语义为「未绑定应用」——此时品牌商账号的应用下拉退化为不限制(见 brandAppNames)。
|
|
// oemfactoryid 是后加字段(这款产品默认由哪家代工厂生产,oem_factory.id)。
|
|
// ⚠️ 不叫 factoryid:DBProduct 里那个 factoryid 是老客户端兼容字段(gorm:"-",值同 brandid),同名会打架。
|
|
if res := adb.Exec("ALTER TABLE " + comm.TableProduct + " ADD COLUMN IF NOT EXISTS oemfactoryid bigint DEFAULT 0"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
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.TableFactoryPublicCode, &pb.DBFactoryPublicCode{}); err != nil {
|
|
return err
|
|
}
|
|
// 设备 MAC 表:**全局一张**,不再按产品切分表。
|
|
// 设备的身份是它自己的 MAC(主键 code = devicemac),productid 只是归属列。
|
|
// 分表时代的教训:pid 既是查询键又要客户端猜,猜错就报「这台 MAC 没登记在该产品下」,
|
|
// 指向数据缺失、实际是 pid 错了,极难排查。详见 migrateLicenseToDeviceMac。
|
|
if err := adb.CreateTable(comm.TableDeviceMac, &pb.DBAuthCode{}); err != nil {
|
|
return err
|
|
}
|
|
var products []*pb.DBProduct
|
|
if err := adb.Find(comm.TableProduct, &products, ""); err != nil && err != postgres.ErrNoDocuments {
|
|
return err
|
|
}
|
|
// 渠道商维度是后加字段:CreateTable 对已存在的表跳过 AutoMigrate,故批次台账与设备表
|
|
// 都要显式补列(幂等)。历史行 channelid 留空,语义为「未分配渠道」。
|
|
for _, t := range []string{comm.TableProductionBatch, comm.TableDeviceMac} {
|
|
if res := adb.Exec("ALTER TABLE " + t + " ADD COLUMN IF NOT EXISTS channelid varchar(16) DEFAULT ''"); res.Error != nil {
|
|
return res.Error
|
|
}
|
|
}
|
|
return migrateLicenseToDeviceMac(products)
|
|
}
|
|
|
|
// migrateLicenseToDeviceMac 把按产品切的 license_<pid十六进制> 分表**合并**进全局
|
|
// device_mac 表。幂等,每次启动跑,跑完即空转。
|
|
//
|
|
// 为什么要合:设备的身份是它自己的 MAC,本来就全局唯一。分表把 pid 变成了查询键,
|
|
// 于是客户端得先猜出 pid 才能校验一台设备——猜错就报「这台 MAC 没登记在该产品下」,
|
|
// 提示词指向数据缺失、实际是 pid 错了。2026-09-04 就踩了这个:一台已正确导入的耳机
|
|
// 因为所属产品的 devicetype 配成了 2(BLE双端)而不是 1(经典蓝牙),被客户端的
|
|
// 「只在 devicetype==1 的产品里挑」过滤掉,兜底退回了另一个产品的 pid,
|
|
// 服务端拿着错 pid 去错的分表里查,自然查不到。
|
|
//
|
|
// 迁移步骤(每张 license_* 表,含没有对应产品的孤儿表——那里面也是真实发出去的设备):
|
|
// ① code := devicemac —— 老数据里 code 是 20 字符的 license 串、devicemac 另存,
|
|
// 两个都指向同一台设备。授权码那套已经删了(utils/license 包不复存在)。
|
|
// ② 按两表**列名交集**做 INSERT ... SELECT ... ON CONFLICT DO NOTHING。
|
|
// 不能 SELECT *:各分表建于不同时期,列不一致(有的多 factoryid/probatch 两列)。
|
|
// ③ productid 为空/0 的行,用表名里的 pid 补上——它是绑定礼、埋点、结算的归属依据。
|
|
//
|
|
// ⚠️ **搬完不删旧表**。这是不可逆数据(已经发到用户手上的设备),留着旧表是出问题时
|
|
// 唯一的回退依据。确认新表跑稳之后再由人手工 DROP。
|
|
func migrateLicenseToDeviceMac(products []*pb.DBProduct) error {
|
|
adb := postgres.GetSys()
|
|
|
|
// device_mac 的列集合,用于和各源表求交集。
|
|
dstCols, err := tableColumns(adb, comm.TableDeviceMac)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(dstCols) == 0 {
|
|
return fmt.Errorf("console.registry: %s 建表失败或无列", comm.TableDeviceMac)
|
|
}
|
|
|
|
var names []string
|
|
if err := adb.Raw(
|
|
"SELECT tablename FROM pg_tables WHERE schemaname = current_schema() AND tablename LIKE ?",
|
|
comm.TableLicense+`\_%`).Scan(&names).Error; err != nil {
|
|
return err
|
|
}
|
|
for _, t := range names {
|
|
// 表名形如 license_b023,后缀是 pid 的十六进制。解不出就当 0,
|
|
// 那种表的行必须自带 productid,否则归属列会是空——只告警不中断。
|
|
pid := pidFromLicenseTable(t)
|
|
|
|
// ① 存量 code 改写成 MAC。WHERE 带 code<>devicemac,跑过一遍后就是 0 行。
|
|
if res := adb.Exec("UPDATE " + t + " SET code = devicemac WHERE devicemac <> '' AND code <> devicemac"); res.Error != nil {
|
|
return res.Error
|
|
} else if res.RowsAffected > 0 {
|
|
log.Infof("console.registry: %s 有 %d 条存量授权码已改写成 MAC", t, res.RowsAffected)
|
|
}
|
|
|
|
srcCols, err := tableColumns(adb, t)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// ② 列名交集。productid 单独处理(要用表名兜底),其余按名字直搬。
|
|
cols := make([]string, 0, len(srcCols))
|
|
for c := range srcCols {
|
|
if c == "productid" {
|
|
continue // 单独处理,见 ③
|
|
}
|
|
if dstCols[c] {
|
|
cols = append(cols, c)
|
|
}
|
|
}
|
|
if len(cols) == 0 {
|
|
continue
|
|
}
|
|
sort.Strings(cols) // map 迭代无序,排一下让生成的 SQL 与日志可复现
|
|
// ③ productid:源表有这一列就用它、为空则回落到表名里的 pid;没这一列就直接用表名。
|
|
srcPid := fmt.Sprintf("%d", pid)
|
|
if srcCols["productid"] {
|
|
srcPid = fmt.Sprintf("COALESCE(NULLIF(productid, 0), %d)", pid)
|
|
}
|
|
sql := fmt.Sprintf(
|
|
"INSERT INTO %s (productid, %s) SELECT %s, %s FROM %s WHERE code <> '' ON CONFLICT (code) DO NOTHING",
|
|
comm.TableDeviceMac, strings.Join(cols, ", "), srcPid, strings.Join(cols, ", "), t)
|
|
res := adb.Exec(sql)
|
|
if res.Error != nil {
|
|
return fmt.Errorf("console.registry: 迁移 %s -> %s 失败: %w", t, comm.TableDeviceMac, res.Error)
|
|
}
|
|
if res.RowsAffected > 0 {
|
|
log.Infof("console.registry: 已把 %s 的 %d 行合并进 %s(旧表保留)", t, res.RowsAffected, comm.TableDeviceMac)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// tableColumns 取一张表的列名集合。表不存在返回空 map(不报错)。
|
|
func tableColumns(adb postgres.ISys, table string) (map[string]bool, error) {
|
|
var cols []string
|
|
if err := adb.Raw(
|
|
"SELECT column_name FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = ?",
|
|
table).Scan(&cols).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
out := make(map[string]bool, len(cols))
|
|
for _, c := range cols {
|
|
out[c] = true
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// pidFromLicenseTable 从 license_<pid十六进制> 解出 pid。解不出返回 0。
|
|
func pidFromLicenseTable(table string) uint32 {
|
|
suffix := strings.TrimPrefix(table, comm.TableLicense+"_")
|
|
if suffix == table || suffix == "" {
|
|
return 0
|
|
}
|
|
n, err := strconv.ParseUint(suffix, 16, 32)
|
|
if err != nil {
|
|
log.Warn("console.registry: 授权码表名解不出 pid,归属列将依赖表内数据",
|
|
log.Field{Key: "table", Value: table})
|
|
return 0
|
|
}
|
|
return uint32(n)
|
|
}
|
|
|