116 lines
4.6 KiB
Python
116 lines
4.6 KiB
Python
#!/usr/bin/env python3
|
||
"""建表脚本:从各模块 models/*.json 生成 DDL 并在 pipeline 库执行。
|
||
|
||
对每个模块,调用 json2ddl 生成 DDL,按分号分割后逐条执行(幂等:DDL 含 DROP TABLE IF EXISTS)。
|
||
|
||
用法:
|
||
py3/bin/python scripts/create_tables.py
|
||
"""
|
||
import sys, os, re, asyncio, subprocess
|
||
|
||
SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
|
||
ROOT_DIR = os.path.dirname(SCRIPT_DIR)
|
||
sys.path.insert(0, os.path.join(ROOT_DIR, 'py3', 'lib', 'python3.10', 'site-packages'))
|
||
sys.path.insert(0, ROOT_DIR)
|
||
|
||
from sqlor.dbpools import DBPools
|
||
from appPublic.jsonConfig import getConfig
|
||
from appPublic.folderUtils import ProgramPath
|
||
from ahserver.serverenv import ServerEnv
|
||
from ahserver.globalEnv import initEnv
|
||
|
||
TABLE_MODULES = ['product_management', 'discount', 'pricing', 'unipay', 'smssend']
|
||
# 幂等建表模块(不 DROP,用 IF NOT EXISTS,避免清空运行时数据如 llm/pipelines/skill_proposals/tasks 等)
|
||
IDEMPOTENT_MODULES = ['pipeline_core', 'pipeline-service']
|
||
|
||
|
||
# appbase 系统级配置表(幂等建表 + 默认参数,不 DROP 以免清空运行时数据)
|
||
_PARAMS_DDL = """
|
||
CREATE TABLE IF NOT EXISTS params (
|
||
`id` VARCHAR(32) comment 'id',
|
||
`params_name` VARCHAR(255) comment '参数名称',
|
||
`params_value` VARCHAR(4000) comment '参数值',
|
||
primary key(id)
|
||
) CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci engine=innodb comment '系统参数表';
|
||
"""
|
||
|
||
_PARAMS_INIT = [
|
||
("workspace_base", "/d/pipeline/workspaces"),
|
||
]
|
||
|
||
|
||
def split_ddl(ddl):
|
||
"""按分号分割 DDL,跳过注释行和空语句。"""
|
||
stmts = []
|
||
for raw in ddl.split(';'):
|
||
lines = [l for l in raw.split('\n') if not l.strip().startswith('--')]
|
||
s = '\n'.join(lines).strip()
|
||
if s:
|
||
stmts.append(s)
|
||
return stmts
|
||
|
||
|
||
async def main():
|
||
config = getConfig(ROOT_DIR, NS={'workdir': ROOT_DIR, 'ProgramPath': ProgramPath()})
|
||
DBPools(config.databases)
|
||
initEnv()
|
||
env = ServerEnv()
|
||
env.get_module_dbname = lambda m: 'pipeline' if 'pipeline' in m else 'sage'
|
||
|
||
json2ddl = os.path.join(ROOT_DIR, 'py3', 'bin', 'json2ddl')
|
||
async with DBPools().sqlorContext('pipeline') as sor:
|
||
for mod in TABLE_MODULES:
|
||
models_dir = os.path.join(ROOT_DIR, 'pkgs', mod, 'models')
|
||
if not os.path.isdir(models_dir):
|
||
print(f' skip {mod}: no models dir')
|
||
continue
|
||
r = subprocess.run([json2ddl, 'mysql', models_dir], capture_output=True, text=True)
|
||
if r.returncode != 0 or not r.stdout.strip():
|
||
print(f' skip {mod}: json2ddl failed')
|
||
continue
|
||
n = 0
|
||
for stmt in split_ddl(r.stdout):
|
||
await sor.execute(stmt, {})
|
||
n += 1
|
||
print(f' tables: {mod} created ({n} statements)')
|
||
|
||
# 幂等建表模块(去 DROP + IF NOT EXISTS,不清理运行时数据)
|
||
for mod in IDEMPOTENT_MODULES:
|
||
models_dir = os.path.join(ROOT_DIR, 'pkgs', mod, 'models')
|
||
if not os.path.isdir(models_dir):
|
||
print(f' skip {mod}: no models dir')
|
||
continue
|
||
r = subprocess.run([json2ddl, 'mysql', models_dir], capture_output=True, text=True)
|
||
if r.returncode != 0 or not r.stdout.strip():
|
||
print(f' skip {mod}: json2ddl failed')
|
||
continue
|
||
n = 0
|
||
for stmt in split_ddl(r.stdout):
|
||
s = re.sub(r'(?i)drop\s+table\s+if\s+exists\s+[\w`]+\s*;?', '', stmt)
|
||
s = re.sub(r'(?i)CREATE\s+TABLE\s+(`?\w+`?)', r'CREATE TABLE IF NOT EXISTS \1', s)
|
||
if not s.strip():
|
||
continue
|
||
# 跳过索引/约束语句(表已存在时索引也已存在,重复建会 Duplicate key name)
|
||
if re.match(r'(?i)\s*(CREATE\s+(UNIQUE\s+)?INDEX|ALTER\s+TABLE)', s):
|
||
continue
|
||
await sor.execute(s, {})
|
||
n += 1
|
||
print(f' tables: {mod} ensured ({n} statements, idempotent)')
|
||
|
||
# appbase params 表幂等建表 + 默认参数
|
||
try:
|
||
await sor.execute(_PARAMS_DDL, {})
|
||
for pname, pval in _PARAMS_INIT:
|
||
await sor.execute(
|
||
"INSERT INTO params (id, params_name, params_value) "
|
||
"VALUES (${id}$, ${n}$, ${v}$) "
|
||
"ON DUPLICATE KEY UPDATE params_value=${v}$",
|
||
{"id": pname, "n": pname, "v": pval})
|
||
print(' tables: appbase params ensured (workspace_base)')
|
||
except Exception as e:
|
||
print(f' WARN: appbase params ensure failed: {e}')
|
||
|
||
|
||
if __name__ == '__main__':
|
||
asyncio.run(main())
|