334 lines
15 KiB
Python
334 lines
15 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"),
|
||
("task_max_retry", "3"),
|
||
]
|
||
|
||
# sd_features 功能/需求表 DDL(手写,避开 json2ddl 的 DEFAULT ''N'' bug;索引内联避免重复建 Duplicate key name)
|
||
_SD_FEATURES_DDL = """
|
||
CREATE TABLE IF NOT EXISTS sd_features (
|
||
`id` VARCHAR(32) NOT NULL comment '主键ID',
|
||
`project_id` VARCHAR(32) NOT NULL comment '项目ID',
|
||
`iteration_id` VARCHAR(32) comment '迭代ID',
|
||
`feature_name` VARCHAR(200) NOT NULL comment '功能名称',
|
||
`description` text comment '功能描述',
|
||
`feature_type` VARCHAR(20) NOT NULL DEFAULT 'new_feature' comment '功能类型',
|
||
`priority` VARCHAR(10) NOT NULL DEFAULT 'P2' comment '优先级',
|
||
`status` VARCHAR(20) NOT NULL DEFAULT 'proposed' comment '功能状态',
|
||
`acceptance_criteria` text comment '验收标准',
|
||
`task_id` VARCHAR(32) comment '关联Pipeline任务ID',
|
||
`created_by` VARCHAR(32) comment '创建人',
|
||
`created_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL comment '创建时间',
|
||
`updated_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP comment '更新时间',
|
||
primary key(id),
|
||
KEY idx_sd_features_project(project_id),
|
||
KEY idx_sd_features_iteration(iteration_id),
|
||
KEY idx_sd_features_status(status)
|
||
) CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci engine=innodb comment '功能/需求表';
|
||
"""
|
||
|
||
# pipeline-service v2 运行时表(原 init.py::_init_v2_tables 启动时建,现迁到部署期)。
|
||
# 手写幂等 DDL:避开 json2ddl 的 DEFAULT ''N'' bug;索引内联避免重复建 Duplicate key name。
|
||
# 铁律:运行期不做任何 schema 变更,DB 全部变化必须在部署期(build.sh → create_tables.py)完成。
|
||
_V2_RUNTIME_TABLES = [
|
||
"""
|
||
CREATE TABLE IF NOT EXISTS pipeline_user_memory (
|
||
`id` varchar(32) NOT NULL,
|
||
`memory_key` varchar(64) NOT NULL,
|
||
`content` text NOT NULL,
|
||
`category` varchar(32) NOT NULL DEFAULT 'memory',
|
||
`scope` varchar(32) NOT NULL DEFAULT 'global',
|
||
`scope_id` varchar(64) NOT NULL DEFAULT '',
|
||
`priority` int(11) NOT NULL DEFAULT 0,
|
||
`access_count` int(11) NOT NULL DEFAULT 0,
|
||
`created_at` datetime NOT NULL DEFAULT current_timestamp(),
|
||
`updated_at` datetime NOT NULL DEFAULT current_timestamp() ON UPDATE current_timestamp(),
|
||
PRIMARY KEY (`id`),
|
||
UNIQUE KEY `uk_memory` (`memory_key`,`category`),
|
||
KEY `idx_category_priority` (`category`,`priority`)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci;
|
||
""",
|
||
"""
|
||
CREATE TABLE IF NOT EXISTS sd_deploy_accounts (
|
||
`id` varchar(32) NOT NULL,
|
||
`user_id` varchar(32) NOT NULL,
|
||
`account_name` varchar(64) NOT NULL,
|
||
`deploy_dir` varchar(500) NOT NULL,
|
||
`status` varchar(20) NOT NULL DEFAULT 'active',
|
||
`sandbox_config` text,
|
||
`created_at` datetime NOT NULL DEFAULT current_timestamp(),
|
||
`updated_at` datetime NOT NULL DEFAULT current_timestamp() ON UPDATE current_timestamp(),
|
||
PRIMARY KEY (`id`),
|
||
UNIQUE KEY `uk_account_name` (`account_name`),
|
||
UNIQUE KEY `uk_user_id` (`user_id`),
|
||
KEY `idx_status` (`status`)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||
""",
|
||
"""
|
||
CREATE TABLE IF NOT EXISTS sd_work_envs (
|
||
`id` varchar(32) NOT NULL,
|
||
`owner_type` varchar(10) NOT NULL,
|
||
`owner_id` varchar(32) NOT NULL,
|
||
`mode` varchar(10) NOT NULL DEFAULT 'local',
|
||
`remote_host` varchar(200),
|
||
`remote_port` int(11) DEFAULT 22,
|
||
`remote_user` varchar(100),
|
||
`remote_key_path` varchar(500),
|
||
`remote_dir` varchar(500),
|
||
`created_at` datetime NOT NULL DEFAULT current_timestamp(),
|
||
`updated_at` datetime NOT NULL DEFAULT current_timestamp() ON UPDATE current_timestamp(),
|
||
PRIMARY KEY (`id`),
|
||
UNIQUE KEY `uk_owner` (`owner_type`,`owner_id`)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||
""",
|
||
]
|
||
|
||
# v2 引擎列迁移(原 init.py::_init_v2_tables 启动时 ALTER 迁移,现迁到部署期,information_schema 幂等)
|
||
_V2_COLUMNS = {
|
||
"pipelines": {"agent_config": "text"},
|
||
"sd_org_settings": {"agent_config": "text"},
|
||
"pipeline_tasks": {"retry_count": "int NOT NULL DEFAULT 0", "last_error": "text",
|
||
"parent_id": "varchar(32) NULL", "depends_on": "text NULL"},
|
||
}
|
||
|
||
# app_audit 审计日志表(原 app_audit/init.py::_init_audit 启动时建,现迁到部署期)
|
||
_AUDIT_LOGS_DDL = """
|
||
CREATE TABLE IF NOT EXISTS sd_audit_logs (
|
||
`id` varchar(32) NOT NULL,
|
||
`user_id` varchar(32),
|
||
`username` varchar(100),
|
||
`action` varchar(50) NOT NULL,
|
||
`target` varchar(200),
|
||
`detail` text,
|
||
`result` varchar(10),
|
||
`client_ip` varchar(64),
|
||
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||
PRIMARY KEY (`id`),
|
||
KEY `idx_user` (`user_id`),
|
||
KEY `idx_action` (`action`),
|
||
KEY `idx_created` (`created_at`)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||
"""
|
||
|
||
# git 级串行锁表:多 worker(含跨主机)对同一 repo 的 git 操作互斥,带 TTL 自动释放
|
||
# (expires_at 过期可被下一个申请者原子接管)。lock_key = sha256(repo_abs_path)。
|
||
_GIT_LOCKS_DDL = """
|
||
CREATE TABLE IF NOT EXISTS pipeline_git_locks (
|
||
`lock_key` varchar(64) NOT NULL,
|
||
`token` varchar(64) NOT NULL,
|
||
`expires_at` datetime NOT NULL,
|
||
PRIMARY KEY (`lock_key`),
|
||
KEY `idx_expires` (`expires_at`)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||
"""
|
||
|
||
|
||
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}')
|
||
|
||
# pipeline_agent_questions 团队沟通列幂等迁移(显式路由版:current_handler 两列替代 escalation 路径快照)
|
||
try:
|
||
_add_cols = {
|
||
"from_agentid": "VARCHAR(64) NULL COMMENT '提问Agent标识'",
|
||
"problem_type": "VARCHAR(50) NULL COMMENT '问题类型'",
|
||
"current_handler_role": "VARCHAR(32) NULL COMMENT '当前处理角色'",
|
||
"current_handler_agentid": "VARCHAR(64) NULL COMMENT '当前处理Agent标识(空=该角色任意)'",
|
||
}
|
||
for col, ddl in _add_cols.items():
|
||
r = await sor.sqlExe(
|
||
"SELECT COUNT(*) as c FROM information_schema.COLUMNS "
|
||
"WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='pipeline_agent_questions' "
|
||
"AND COLUMN_NAME=${col}$", {"col": col})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not r or getattr(r[0], 'c', 0) == 0:
|
||
await sor.execute(
|
||
f"ALTER TABLE pipeline_agent_questions ADD COLUMN {col} {ddl}", {})
|
||
# 废弃列清理(早期路径快照方案,已改为显式 handler)
|
||
for col in ["escalation_path", "escalation_pos"]:
|
||
r = await sor.sqlExe(
|
||
"SELECT COUNT(*) as c FROM information_schema.COLUMNS "
|
||
"WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='pipeline_agent_questions' "
|
||
"AND COLUMN_NAME=${col}$", {"col": col})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if r and getattr(r[0], 'c', 0) > 0:
|
||
await sor.execute(
|
||
f"ALTER TABLE pipeline_agent_questions DROP COLUMN {col}", {})
|
||
print(' migrate: pipeline_agent_questions team-communication columns ensured')
|
||
except Exception as e:
|
||
print(f' WARN: pipeline_agent_questions migrate failed: {e}')
|
||
|
||
# sd_features 功能/需求表(pipeline-sdlc 业务表)幂等建表——手写 DDL 避开 json2ddl 的 DEFAULT ''N'' bug
|
||
try:
|
||
await sor.execute(_SD_FEATURES_DDL, {})
|
||
print(' tables: sd_features ensured (idempotent)')
|
||
except Exception as e:
|
||
print(f' WARN: sd_features ensure failed: {e}')
|
||
|
||
# pipeline-service v2 运行时表(原启动时建,现迁到部署期幂等建表)
|
||
try:
|
||
for ddl in _V2_RUNTIME_TABLES:
|
||
await sor.execute(ddl, {})
|
||
print(' tables: pipeline-service v2 runtime tables ensured (pipeline_user_memory/sd_deploy_accounts/sd_work_envs)')
|
||
except Exception as e:
|
||
print(f' WARN: v2 runtime tables ensure failed: {e}')
|
||
|
||
# v2 引擎列迁移(原启动时 ALTER,现迁到部署期,information_schema 查列存在性幂等)
|
||
try:
|
||
for tbl, cols in _V2_COLUMNS.items():
|
||
for col, ddl in cols.items():
|
||
r = await sor.sqlExe(
|
||
"SELECT COUNT(*) as c FROM information_schema.COLUMNS "
|
||
"WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=${tbl}$ "
|
||
"AND COLUMN_NAME=${col}$", {"tbl": tbl, "col": col})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not r or getattr(r[0], 'c', 0) == 0:
|
||
await sor.execute(f"ALTER TABLE {tbl} ADD COLUMN {col} {ddl}", {})
|
||
print(' migrate: v2 engine columns ensured (agent_config/retry_count/last_error)')
|
||
except Exception as e:
|
||
print(f' WARN: v2 engine column migrate failed: {e}')
|
||
|
||
# app_audit 审计日志表(原启动时建,现迁到部署期幂等建表)
|
||
try:
|
||
await sor.execute(_AUDIT_LOGS_DDL, {})
|
||
print(' tables: sd_audit_logs ensured (idempotent)')
|
||
except Exception as e:
|
||
print(f' WARN: sd_audit_logs ensure failed: {e}')
|
||
|
||
# git 级串行锁表(跨主机多 worker 对同 repo git 操作互斥)
|
||
try:
|
||
await sor.execute(_GIT_LOCKS_DDL, {})
|
||
print(' tables: pipeline_git_locks ensured (idempotent)')
|
||
except Exception as e:
|
||
print(f' WARN: pipeline_git_locks ensure failed: {e}')
|
||
|
||
# sd_iterations 迭代序号 seq_no 列 + 存量回填 + 脏状态修复(active→in_progress)
|
||
# 「当前迭代」改为显式 status='in_progress' 标记,迭代按 seq_no 顺序编号推进,
|
||
# 不再靠 created_at 推断当前迭代。
|
||
try:
|
||
r = await sor.sqlExe(
|
||
"SELECT COUNT(*) as c FROM information_schema.COLUMNS "
|
||
"WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME='sd_iterations' "
|
||
"AND COLUMN_NAME='seq_no'", {})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if not r or getattr(r[0], 'c', 0) == 0:
|
||
await sor.execute("ALTER TABLE sd_iterations ADD COLUMN seq_no INT NOT NULL DEFAULT 0", {})
|
||
# 存量回填:seq_no=0 的迭代按每项目 created_at ASC 顺序赋 1,2,3...
|
||
await sor.sqlExe(
|
||
"UPDATE sd_iterations t JOIN ("
|
||
"SELECT id, ROW_NUMBER() OVER (PARTITION BY project_id ORDER BY created_at ASC, id ASC) AS rn "
|
||
"FROM sd_iterations WHERE seq_no = 0) x ON x.id = t.id "
|
||
"SET t.seq_no = x.rn", {})
|
||
# 脏状态修复:迭代状态合法值 planning/in_progress/completed/cancelled,无 active
|
||
await sor.sqlExe(
|
||
"UPDATE sd_iterations SET status='in_progress', updated_at=NOW() WHERE status='active'", {})
|
||
print(' migrate: sd_iterations seq_no column + backfill + active→in_progress ensured')
|
||
except Exception as e:
|
||
print(f' WARN: sd_iterations migrate failed: {e}')
|
||
|
||
|
||
if __name__ == '__main__':
|
||
asyncio.run(main())
|