From ff7f03e7ffe833c02d7229dcc8711375059ea1be Mon Sep 17 00:00:00 2001 From: yumoqing Date: Sat, 4 Jul 2026 11:04:22 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20FOMS=20v1.0=20-=20Failover=20Management?= =?UTF-8?q?=20System=20(MySQL=E4=B8=BB=E4=BB=8E=E5=A4=8D=E5=88=B6+?= =?UTF-8?q?=E6=96=87=E4=BB=B6=E5=90=8C=E6=AD=A5+=E6=95=85=E9=9A=9C?= =?UTF-8?q?=E5=88=87=E6=8D=A2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 20 ++ README.md | 61 ++++ app/foms.py | 32 ++ app/global_func.py | 12 + build.sh | 51 ++++ conf/config.json | 42 +++ foms/__init__.py | 8 + foms/init.py | 366 +++++++++++++++++++++++ init/data.yaml | 134 +++++++++ json/clusters_list.json | 22 ++ json/nodes_list.json | 21 ++ json/sync_databases_list.json | 22 ++ json/sync_directories_list.json | 22 ++ models/clusters.json | 30 ++ models/nodes.json | 31 ++ models/switchover_logs.json | 34 +++ models/sync_databases.json | 35 +++ models/sync_directories.json | 33 ++ pyproject.toml | 17 ++ scripts/foms-agent.py | 260 ++++++++++++++++ scripts/load_path.py | 84 ++++++ start.sh | 6 + stop.sh | 9 + wwwroot/api/agent_heartbeat.dspy | 7 + wwwroot/api/agent_report_health.dspy | 3 + wwwroot/api/clusters_create.dspy | 10 + wwwroot/api/clusters_delete.dspy | 10 + wwwroot/api/clusters_update.dspy | 10 + wwwroot/api/dashboard_stats.dspy | 3 + wwwroot/api/nodes_create.dspy | 10 + wwwroot/api/nodes_delete.dspy | 10 + wwwroot/api/nodes_update.dspy | 10 + wwwroot/api/sync_databases_create.dspy | 10 + wwwroot/api/sync_databases_delete.dspy | 10 + wwwroot/api/sync_databases_update.dspy | 10 + wwwroot/api/sync_directories_create.dspy | 10 + wwwroot/api/sync_directories_delete.dspy | 10 + wwwroot/api/sync_directories_update.dspy | 10 + wwwroot/api/trigger_switchback.dspy | 6 + wwwroot/api/trigger_switchover.dspy | 6 + wwwroot/index.ui | 87 ++++++ 41 files changed, 1584 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 app/foms.py create mode 100644 app/global_func.py create mode 100755 build.sh create mode 100644 conf/config.json create mode 100644 foms/__init__.py create mode 100644 foms/init.py create mode 100644 init/data.yaml create mode 100644 json/clusters_list.json create mode 100644 json/nodes_list.json create mode 100644 json/sync_databases_list.json create mode 100644 json/sync_directories_list.json create mode 100644 models/clusters.json create mode 100644 models/nodes.json create mode 100644 models/switchover_logs.json create mode 100644 models/sync_databases.json create mode 100644 models/sync_directories.json create mode 100644 pyproject.toml create mode 100644 scripts/foms-agent.py create mode 100644 scripts/load_path.py create mode 100755 start.sh create mode 100755 stop.sh create mode 100644 wwwroot/api/agent_heartbeat.dspy create mode 100644 wwwroot/api/agent_report_health.dspy create mode 100644 wwwroot/api/clusters_create.dspy create mode 100644 wwwroot/api/clusters_delete.dspy create mode 100644 wwwroot/api/clusters_update.dspy create mode 100644 wwwroot/api/dashboard_stats.dspy create mode 100644 wwwroot/api/nodes_create.dspy create mode 100644 wwwroot/api/nodes_delete.dspy create mode 100644 wwwroot/api/nodes_update.dspy create mode 100644 wwwroot/api/sync_databases_create.dspy create mode 100644 wwwroot/api/sync_databases_delete.dspy create mode 100644 wwwroot/api/sync_databases_update.dspy create mode 100644 wwwroot/api/sync_directories_create.dspy create mode 100644 wwwroot/api/sync_directories_delete.dspy create mode 100644 wwwroot/api/sync_directories_update.dspy create mode 100644 wwwroot/api/trigger_switchback.dspy create mode 100644 wwwroot/api/trigger_switchover.dspy create mode 100644 wwwroot/index.ui diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..05e4cdb --- /dev/null +++ b/.gitignore @@ -0,0 +1,20 @@ +__pycache__/ +*.pyc +*.pyo +*.egg-info/ +build/ +dist/ +*.pid +logs/ +files/ +pkgs/ +bricks/ +py3/ +.env +*.swp +*.swo +# CRUD-generated directories (xls2ui) +wwwroot/clusters_list/ +wwwroot/nodes_list/ +wwwroot/sync_databases_list/ +wwwroot/sync_directories_list/ diff --git a/README.md b/README.md new file mode 100644 index 0000000..6cdeb34 --- /dev/null +++ b/README.md @@ -0,0 +1,61 @@ +# FOMS — Failover Management System + +主备机切换管理系统。运行在管理中心 C 上,管理主机 A 和备机 B 之间的: + +- **MySQL 主从复制** — 实时数据同步 +- **文件目录同步** — rsync 定期同步 +- **健康监控** — 节点心跳 + 复制状态 +- **故障切换** — 自动检测 A 故障 → 切换到 B +- **故障切回** — A 修复后 → 切回 A + +## 架构 + +``` +管理中心 C (FOMS Web App) + │ + ├── 主机 A ←──MySQL复制/rsync──→ 备机 B + │ + └── Agent 守护进程 (部署在 A & B) +``` + +## 快速开始 + +```bash +# 1. 构建 +./build.sh + +# 2. 配置 .env +vim .env + +# 3. 初始化数据库 (手动执行 models/mysql.ddl.sql) +mysql -u root < models/mysql.ddl.sql + +# 4. 注册RBAC权限 +cd py3/bin && python ../../scripts/load_path.py + +# 5. 启动 +./start.sh +``` + +## 模块 + +| 模块 | 路径 | 说明 | +|------|------|------| +| 集群管理 | /clusters_list | 配置主备集群 | +| 节点管理 | /nodes_list | 注册A/B节点(SSH信息) | +| 数据库同步 | /sync_databases_list | MySQL主从复制配置 | +| 文件同步 | /sync_directories_list | rsync目录同步配置 | + +## Agent部署 + +在主机A和备机B上: + +```bash +# 配置环境变量 +export FOMS_MANAGEMENT_URL=http://C的IP:9080 +export FOMS_NODE_ID=<节点ID> +export FOMS_AGENT_TOKEN=<安全令牌> + +# 运行Agent +python3 scripts/foms-agent.py +``` diff --git a/app/foms.py b/app/foms.py new file mode 100644 index 0000000..4c93f95 --- /dev/null +++ b/app/foms.py @@ -0,0 +1,32 @@ +#!/usr/bin/env python3 +"""FOMS — Failover Management System 主入口""" +import os, sys + +# Add root so sibling modules are importable +root_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +sys.path.insert(0, root_dir) + +from bricks_for_python.init import load_pybricks +from ahserver.webapp import webapp +from ahserver.serverenv import ServerEnv + +# Foundation modules +from rbac.init import load_rbac +from appbase.init import load_appbase + +# Business modules +from foms.init import load_foms + +from app.global_func import set_globalvariable + + +def init(): + set_globalvariable() + load_pybricks() + load_appbase() + load_rbac() + load_foms() + + +if __name__ == '__main__': + webapp(init) diff --git a/app/global_func.py b/app/global_func.py new file mode 100644 index 0000000..02859dd --- /dev/null +++ b/app/global_func.py @@ -0,0 +1,12 @@ +"""FOMS global functions — injected into dspy context""" +from ahserver.serverenv import ServerEnv + + +def get_module_dbname(mname): + """All FOMS modules share one database.""" + return 'foms' + + +def set_globalvariable(): + env = ServerEnv() + env.get_module_dbname = get_module_dbname diff --git a/build.sh b/build.sh new file mode 100755 index 0000000..aff3149 --- /dev/null +++ b/build.sh @@ -0,0 +1,51 @@ +#!/usr/bin/env bash +set -e +cdir=$(pwd) +echo "=== FOMS Build ===" + +# 1. Create venv +[ ! -d py3 ] && python3 -m venv py3 +source py3/bin/activate + +# 2. Install foundation packages +mkdir -p pkgs +for m in apppublic sqlor ahserver bricks-for-python xls2ddl; do + cd $cdir/pkgs + [ ! -d "$m" ] && git clone https://git.opencomputing.cn/yumoqing/$m + cd $m && $cdir/py3/bin/pip install -e . +done + +# 3. Build bricks frontend +cd $cdir/pkgs +[ ! -d "bricks" ] && git clone https://git.opencomputing.cn/yumoqing/bricks +cd bricks/bricks && ./build.sh +ln -sf $cdir/pkgs/bricks/dist $cdir/bricks + +# 4. Install auth modules +for m in appbase rbac; do + cd $cdir/pkgs + [ ! -d "$m" ] && git clone https://git.opencomputing.cn/yumoqing/$m + cd $m && $cdir/py3/bin/pip install -e . +done + +# 5. Install FOMS +cd $cdir && $cdir/py3/bin/pip install -e . + +# 6. Generate DDL +cd $cdir +[ -d models ] && json2ddl mysql models/ > $cdir/models/mysql.ddl.sql 2>/dev/null || echo " skip ddl" + +# 7. Generate CRUD UI (xls2ui for json/) +cd $cdir +if [ -d json ]; then + cd json + for f in *.json; do + echo " xls2ui $f" + xls2ui -m ../models -o ../wwwroot foms "$f" 2>/dev/null || echo " skipped" + done +fi + +# 8. Create runtime dirs +mkdir -p $cdir/logs $cdir/files + +echo "=== FOMS Build Complete ===" diff --git a/conf/config.json b/conf/config.json new file mode 100644 index 0000000..a1dfd60 --- /dev/null +++ b/conf/config.json @@ -0,0 +1,42 @@ +{ + "password_key": "foms_secret_key_change_me", + "logger": { + "name": "foms", + "level": "INFO", + "file": "$[workdir]$/logs/foms.log" + }, + "filesroot": "$[workdir]$/files", + "databases": { + "foms": { + "driver": "aiomysql", + "kwargs": { + "host": "127.0.0.1", + "port": 3306, + "user": "root", + "password": "", + "db": "foms", + "charset": "utf8mb4", + "autocommit": true + } + } + }, + "website": { + "paths": [ + ["$[workdir]$/wwwroot", ""], + ["$[workdir]$/bricks", "/bricks"] + ], + "port": 9080, + "indexes": ["index.ui"], + "session_max_time": 86400, + "session_issue_time": 3600, + "processors": [ + [".dspy", "dspy"], + [".ui", "bui"], + [".tmpl", "tmpl"], + [".js", "staticfile"], + [".css", "staticfile"], + [".png", "staticfile"], + [".svg", "staticfile"] + ] + } +} diff --git a/foms/__init__.py b/foms/__init__.py new file mode 100644 index 0000000..d4d6580 --- /dev/null +++ b/foms/__init__.py @@ -0,0 +1,8 @@ +from .init import ( + create_cluster, update_cluster, delete_cluster, get_cluster_list, + create_node, update_node, delete_node, get_node_list, + create_sync_db, update_sync_db, delete_sync_db, get_sync_db_list, + create_sync_dir, update_sync_dir, delete_sync_dir, get_sync_dir_list, + agent_heartbeat, agent_report_health, trigger_switchover, trigger_switchback, + get_dashboard_stats, +) diff --git a/foms/init.py b/foms/init.py new file mode 100644 index 0000000..f6c7900 --- /dev/null +++ b/foms/init.py @@ -0,0 +1,366 @@ +"""FOMS business logic — CRUD operations + failover orchestration""" +import json +from datetime import datetime +from appPublic.uniqueID import getID +from appPublic.timeUtils import curDateString, timestampstr +from sqlor.dbpools import DBPools, get_sor_context +from ahserver.serverenv import ServerEnv + +DBNAME = 'foms' + + +# ═══════════════════════════════════════════ +# Cluster CRUD +# ═══════════════════════════════════════════ + +async def create_cluster(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['id'] = getID() + ns['status'] = 'stopped' + ns['created_at'] = curDateString() + ns['updated_at'] = curDateString() + # strip _text fields from Tabular + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.C('clusters', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '集群创建成功', 'type': 'success'}} + + +async def update_cluster(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.U('clusters', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '集群更新成功', 'type': 'success'}} + + +async def delete_cluster(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + _id = params_kw.get('id', '') if hasattr(params_kw, 'get') else params_kw['id'] + await sor.D('clusters', {'id': _id}) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '集群已删除', 'type': 'success'}} + + +# ═══════════════════════════════════════════ +# Node CRUD +# ═══════════════════════════════════════════ + +async def create_node(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['id'] = getID() + ns['status'] = 'offline' + ns['created_at'] = curDateString() + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.C('nodes', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '节点创建成功', 'type': 'success'}} + + +async def update_node(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.U('nodes', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '节点更新成功', 'type': 'success'}} + + +async def delete_node(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + _id = params_kw.get('id', '') if hasattr(params_kw, 'get') else params_kw['id'] + await sor.D('nodes', {'id': _id}) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '节点已删除', 'type': 'success'}} + + +# ═══════════════════════════════════════════ +# Sync Database CRUD +# ═══════════════════════════════════════════ + +async def create_sync_db(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['id'] = getID() + ns['sync_status'] = 'not_configured' + ns['created_at'] = curDateString() + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.C('sync_databases', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '数据库同步配置创建成功', 'type': 'success'}} + + +async def update_sync_db(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.U('sync_databases', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '数据库同步配置更新成功', 'type': 'success'}} + + +async def delete_sync_db(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + _id = params_kw.get('id', '') if hasattr(params_kw, 'get') else params_kw['id'] + await sor.D('sync_databases', {'id': _id}) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '数据库同步配置已删除', 'type': 'success'}} + + +# ═══════════════════════════════════════════ +# Sync Directory CRUD +# ═══════════════════════════════════════════ + +async def create_sync_dir(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['id'] = getID() + ns['created_at'] = curDateString() + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.C('sync_directories', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '目录同步配置创建成功', 'type': 'success'}} + + +async def update_sync_dir(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + ns['updated_at'] = curDateString() + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + await sor.U('sync_directories', ns) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '目录同步配置更新成功', 'type': 'success'}} + + +async def delete_sync_dir(request, params_kw): + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + _id = params_kw.get('id', '') if hasattr(params_kw, 'get') else params_kw['id'] + await sor.D('sync_directories', {'id': _id}) + return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '目录同步配置已删除', 'type': 'success'}} + + +# ═══════════════════════════════════════════ +# Agent Communication +# ═══════════════════════════════════════════ + +async def agent_heartbeat(request, params_kw): + """Agent心跳上报""" + node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else '' + version = params_kw.get('version', '') if hasattr(params_kw, 'get') else '' + if not node_id: + return {'status': 'error', 'message': 'node_id required'} + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + await sor.U('nodes', { + 'id': node_id, + 'status': 'online', + 'agent_version': version, + 'last_heartbeat': curDateString(), + 'updated_at': curDateString(), + }) + return {'status': 'ok', 'server_time': curDateString()} + + +async def agent_report_health(request, params_kw): + """Agent健康报告(包含同步状态)""" + env = ServerEnv() + node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else '' + async with get_sor_context(env, DBNAME) as sor: + # Update sync_database statuses + db_statuses = params_kw.get('db_statuses', []) if hasattr(params_kw, 'get') else [] + if isinstance(db_statuses, str): + db_statuses = json.loads(db_statuses) + for db_stat in db_statuses: + await sor.U('sync_databases', { + 'id': db_stat.get('id'), + 'sync_status': db_stat.get('status'), + 'seconds_behind': db_stat.get('seconds_behind'), + 'last_error': db_stat.get('error'), + 'updated_at': curDateString(), + }) + # Update sync_directory statuses + dir_statuses = params_kw.get('dir_statuses', []) if hasattr(params_kw, 'get') else [] + if isinstance(dir_statuses, str): + dir_statuses = json.loads(dir_statuses) + for dir_stat in dir_statuses: + await sor.U('sync_directories', { + 'id': dir_stat.get('id'), + 'last_sync_at': curDateString(), + 'last_sync_result': dir_stat.get('result'), + 'updated_at': curDateString(), + }) + return {'status': 'ok'} + + +# ═══════════════════════════════════════════ +# Failover Orchestration +# ═══════════════════════════════════════════ + +async def trigger_switchover(request, params_kw): + """触发故障切换 A→B""" + env = ServerEnv() + cluster_id = params_kw.get('cluster_id', '') if hasattr(params_kw, 'get') else '' + reason = params_kw.get('reason', '手动触发') if hasattr(params_kw, 'get') else '手动触发' + log_id = getID() + now = curDateString() + async with get_sor_context(env, DBNAME) as sor: + # Get cluster info + clusters = await sor.R('clusters', {'id': cluster_id}) + if not clusters: + return {'status': 'error', 'message': '集群不存在'} + cluster = clusters[0] + # Log the switchover + await sor.C('switchover_logs', { + 'id': log_id, + 'cluster_id': cluster_id, + 'event_type': 'failover', + 'from_node_id': getattr(cluster, 'master_node_id', ''), + 'to_node_id': getattr(cluster, 'standby_node_id', ''), + 'trigger': 'manual', + 'reason': reason, + 'status': 'in_progress', + 'started_at': now, + 'created_at': now, + }) + # Update cluster status + await sor.U('clusters', {'id': cluster_id, 'status': 'switching', 'updated_at': now}) + return { + 'status': 'ok', + 'message': '切换流程已触发', + 'log_id': log_id, + } + + +async def trigger_switchback(request, params_kw): + """触发切回 B→A""" + env = ServerEnv() + cluster_id = params_kw.get('cluster_id', '') if hasattr(params_kw, 'get') else '' + log_id = getID() + now = curDateString() + async with get_sor_context(env, DBNAME) as sor: + clusters = await sor.R('clusters', {'id': cluster_id}) + if not clusters: + return {'status': 'error', 'message': '集群不存在'} + cluster = clusters[0] + await sor.C('switchover_logs', { + 'id': log_id, + 'cluster_id': cluster_id, + 'event_type': 'failback', + 'from_node_id': getattr(cluster, 'standby_node_id', ''), + 'to_node_id': getattr(cluster, 'master_node_id', ''), + 'trigger': 'manual', + 'reason': '手动切回', + 'status': 'in_progress', + 'started_at': now, + 'created_at': now, + }) + await sor.U('clusters', {'id': cluster_id, 'status': 'switching', 'updated_at': now}) + return { + 'status': 'ok', + 'message': '切回流程已触发', + 'log_id': log_id, + } + + +# ═══════════════════════════════════════════ +# Dashboard +# ═══════════════════════════════════════════ + +async def get_dashboard_stats(request): + """仪表盘统计数据""" + env = ServerEnv() + async with get_sor_context(env, DBNAME) as sor: + clusters = await sor.sqlExe('SELECT count(*) as cnt FROM clusters', {}) + nodes_online = await sor.sqlExe("SELECT count(*) as cnt FROM nodes WHERE status='online'", {}) + nodes_offline = await sor.sqlExe("SELECT count(*) as cnt FROM nodes WHERE status='offline'", {}) + dbs_synced = await sor.sqlExe("SELECT count(*) as cnt FROM sync_databases WHERE sync_status='running'", {}) + dbs_error = await sor.sqlExe("SELECT count(*) as cnt FROM sync_databases WHERE sync_status IN ('error','failed')", {}) + dirs_synced = await sor.sqlExe("SELECT count(*) as cnt FROM sync_directories WHERE last_sync_result='success'", {}) + recent_logs = await sor.sqlExe( + 'SELECT event_type, status, reason, created_at FROM switchover_logs ORDER BY created_at DESC LIMIT 10', {} + ) + return { + 'total_clusters': clusters[0].cnt if clusters else 0, + 'nodes_online': nodes_online[0].cnt if nodes_online else 0, + 'nodes_offline': nodes_offline[0].cnt if nodes_offline else 0, + 'dbs_synced': dbs_synced[0].cnt if dbs_synced else 0, + 'dbs_error': dbs_error[0].cnt if dbs_error else 0, + 'dirs_synced': dirs_synced[0].cnt if dirs_synced else 0, + 'recent_logs': [{ + 'event_type': getattr(r, 'event_type', ''), + 'status': getattr(r, 'status', ''), + 'reason': getattr(r, 'reason', ''), + 'created_at': str(getattr(r, 'created_at', '')), + } for r in (recent_logs or [])], + } + + +# ═══════════════════════════════════════════ +# ServerEnv Registration +# ═══════════════════════════════════════════ + +def load_foms(): + env = ServerEnv() + # Cluster CRUD + env.create_cluster = create_cluster + env.create_clusters = create_cluster + env.update_cluster = update_cluster + env.update_clusters = update_cluster + env.delete_cluster = delete_cluster + env.delete_clusters = delete_cluster + # Node CRUD + env.create_node = create_node + env.create_nodes = create_node + env.update_node = update_node + env.update_nodes = update_node + env.delete_node = delete_node + env.delete_nodes = delete_node + # Sync DB CRUD + env.create_sync_db = create_sync_db + env.create_sync_databases = create_sync_db + env.update_sync_db = update_sync_db + env.update_sync_databases = update_sync_db + env.delete_sync_db = delete_sync_db + env.delete_sync_databases = delete_sync_db + # Sync Dir CRUD + env.create_sync_dir = create_sync_dir + env.create_sync_directories = create_sync_dir + env.update_sync_dir = update_sync_dir + env.update_sync_directories = update_sync_dir + env.delete_sync_dir = delete_sync_dir + env.delete_sync_directories = delete_sync_dir + # Agent + env.agent_heartbeat = agent_heartbeat + env.agent_report_health = agent_report_health + # Failover + env.trigger_switchover = trigger_switchover + env.trigger_switchback = trigger_switchback + # Dashboard + env.get_dashboard_stats = get_dashboard_stats + return True diff --git a/init/data.yaml b/init/data.yaml new file mode 100644 index 0000000..c5ebaf3 --- /dev/null +++ b/init/data.yaml @@ -0,0 +1,134 @@ +appcodes: +- id: cluster_status + name: 集群状态 + hierarchy_flg: 0 +- id: node_role + name: 节点角色 + hierarchy_flg: 0 +- id: node_status + name: 节点状态 + hierarchy_flg: 0 +- id: sync_status + name: 同步状态 + hierarchy_flg: 0 +- id: sync_result + name: 同步结果 + hierarchy_flg: 0 +- id: switchover_event_type + name: 切换事件类型 + hierarchy_flg: 0 +- id: switchover_trigger + name: 触发方式 + hierarchy_flg: 0 +- id: switchover_status + name: 切换执行状态 + hierarchy_flg: 0 + +appcodes_kv: +- id: cluster_status_running + parentid: cluster_status + k: running + v: 运行中 +- id: cluster_status_stopped + parentid: cluster_status + k: stopped + v: 已停止 +- id: cluster_status_switching + parentid: cluster_status + k: switching + v: 切换中 +- id: cluster_status_error + parentid: cluster_status + k: error + v: 异常 + +- id: node_role_master + parentid: node_role + k: master + v: 主机 +- id: node_role_standby + parentid: node_role + k: standby + v: 备机 + +- id: node_status_online + parentid: node_status + k: online + v: 在线 +- id: node_status_offline + parentid: node_status + k: offline + v: 离线 +- id: node_status_error + parentid: node_status + k: error + v: 异常 + +- id: sync_status_not_configured + parentid: sync_status + k: not_configured + v: 未配置 +- id: sync_status_running + parentid: sync_status + k: running + v: 同步中 +- id: sync_status_error + parentid: sync_status + k: error + v: 异常 +- id: sync_status_failed + parentid: sync_status + k: failed + v: 失败 + +- id: sync_result_success + parentid: sync_result + k: success + v: 成功 +- id: sync_result_failed + parentid: sync_result + k: failed + v: 失败 +- id: sync_result_running + parentid: sync_result + k: running + v: 同步中 + +- id: switchover_event_failover + parentid: switchover_event_type + k: failover + v: 故障切换 +- id: switchover_event_failback + parentid: switchover_event_type + k: failback + v: 故障切回 +- id: switchover_event_manual + parentid: switchover_event_type + k: manual + v: 手动操作 + +- id: switchover_trigger_auto + parentid: switchover_trigger + k: auto + v: 自动检测 +- id: switchover_trigger_manual + parentid: switchover_trigger + k: manual + v: 手动触发 + +- id: switchover_status_pending + parentid: switchover_status + k: pending + v: 等待执行 +- id: switchover_status_in_progress + parentid: switchover_status + k: in_progress + v: 执行中 +- id: switchover_status_success + parentid: switchover_status + k: success + v: 成功 +- id: switchover_status_failed + parentid: switchover_status + k: failed + v: 失败 diff --git a/json/clusters_list.json b/json/clusters_list.json new file mode 100644 index 0000000..c50fa79 --- /dev/null +++ b/json/clusters_list.json @@ -0,0 +1,22 @@ +{ + "tblname": "clusters", + "alias": "clusters_list", + "title": "集群管理", + "params": { + "sortby": ["updated_at desc"], + "browserfields": { + "exclouded": ["id"], + "alters": { + "status": {"uitype": "code"}, + "master_node_id": {"uitype": "code", "valueField": "id", "textField": "host"}, + "standby_node_id": {"uitype": "code", "valueField": "id", "textField": "host"} + } + }, + "editexclouded": ["id", "created_at", "updated_at", "status"], + "editable": { + "new_data_url": "{{entire_url('api/clusters_create.dspy')}}", + "update_data_url": "{{entire_url('api/clusters_update.dspy')}}", + "delete_data_url": "{{entire_url('api/clusters_delete.dspy')}}" + } + } +} diff --git a/json/nodes_list.json b/json/nodes_list.json new file mode 100644 index 0000000..c290678 --- /dev/null +++ b/json/nodes_list.json @@ -0,0 +1,21 @@ +{ + "tblname": "nodes", + "alias": "nodes_list", + "title": "节点管理", + "params": { + "sortby": ["updated_at desc"], + "browserfields": { + "exclouded": ["id"], + "alters": { + "role": {"uitype": "code"}, + "status": {"uitype": "code"} + } + }, + "editexclouded": ["id", "created_at", "updated_at", "status", "agent_version", "last_heartbeat"], + "editable": { + "new_data_url": "{{entire_url('api/nodes_create.dspy')}}", + "update_data_url": "{{entire_url('api/nodes_update.dspy')}}", + "delete_data_url": "{{entire_url('api/nodes_delete.dspy')}}" + } + } +} diff --git a/json/sync_databases_list.json b/json/sync_databases_list.json new file mode 100644 index 0000000..7738b99 --- /dev/null +++ b/json/sync_databases_list.json @@ -0,0 +1,22 @@ +{ + "tblname": "sync_databases", + "alias": "sync_databases_list", + "title": "数据库同步配置", + "params": { + "sortby": ["updated_at desc"], + "browserfields": { + "exclouded": ["id", "repl_password"], + "alters": { + "sync_status": {"uitype": "code"}, + "enabled": {"uitype": "code", "data": [{"value": "1", "text": "启用"}, {"value": "0", "text": "禁用"}]}, + "cluster_id": {"uitype": "code", "valueField": "id", "textField": "name"} + } + }, + "editexclouded": ["id", "created_at", "updated_at", "sync_status", "seconds_behind", "last_error"], + "editable": { + "new_data_url": "{{entire_url('api/sync_databases_create.dspy')}}", + "update_data_url": "{{entire_url('api/sync_databases_update.dspy')}}", + "delete_data_url": "{{entire_url('api/sync_databases_delete.dspy')}}" + } + } +} diff --git a/json/sync_directories_list.json b/json/sync_directories_list.json new file mode 100644 index 0000000..4401694 --- /dev/null +++ b/json/sync_directories_list.json @@ -0,0 +1,22 @@ +{ + "tblname": "sync_directories", + "alias": "sync_directories_list", + "title": "文件目录同步配置", + "params": { + "sortby": ["updated_at desc"], + "browserfields": { + "exclouded": ["id"], + "alters": { + "last_sync_result": {"uitype": "code"}, + "enabled": {"uitype": "code", "data": [{"value": "1", "text": "启用"}, {"value": "0", "text": "禁用"}]}, + "cluster_id": {"uitype": "code", "valueField": "id", "textField": "name"} + } + }, + "editexclouded": ["id", "created_at", "updated_at", "last_sync_at", "last_sync_result"], + "editable": { + "new_data_url": "{{entire_url('api/sync_directories_create.dspy')}}", + "update_data_url": "{{entire_url('api/sync_directories_update.dspy')}}", + "delete_data_url": "{{entire_url('api/sync_directories_delete.dspy')}}" + } + } +} diff --git a/models/clusters.json b/models/clusters.json new file mode 100644 index 0000000..6923ee4 --- /dev/null +++ b/models/clusters.json @@ -0,0 +1,30 @@ +{ + "summary": [ + { + "name": "clusters", + "title": "集群定义", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "集群ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "name", "title": "集群名称", "type": "str", "length": 128, "nullable": "no"}, + {"name": "description", "title": "集群描述", "type": "str", "length": 512}, + {"name": "master_node_id", "title": "主机节点", "type": "str", "length": 32}, + {"name": "standby_node_id", "title": "备机节点", "type": "str", "length": 32}, + {"name": "check_interval", "title": "检测间隔(秒)", "type": "int", "nullable": "no", "default": "10"}, + {"name": "failover_threshold", "title": "故障判定阈值(次)", "type": "int", "nullable": "no", "default": "3"}, + {"name": "status", "title": "集群状态", "type": "str", "length": 32, "nullable": "no", "default": "stopped"}, + {"name": "created_at", "title": "创建时间", "type": "datetime"}, + {"name": "updated_at", "title": "更新时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_clusters_name", "idxtype": "unique", "idxfields": ["name"]} + ], + "codes": [ + {"field": "status", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='cluster_status'"}, + {"field": "master_node_id", "table": "nodes", "valuefield": "id", "textfield": "host"}, + {"field": "standby_node_id", "table": "nodes", "valuefield": "id", "textfield": "host"} + ] +} diff --git a/models/nodes.json b/models/nodes.json new file mode 100644 index 0000000..c87b743 --- /dev/null +++ b/models/nodes.json @@ -0,0 +1,31 @@ +{ + "summary": [ + { + "name": "nodes", + "title": "节点管理", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "节点ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "host", "title": "主机地址", "type": "str", "length": 255, "nullable": "no"}, + {"name": "name", "title": "节点名称", "type": "str", "length": 128}, + {"name": "ip", "title": "IP地址", "type": "str", "length": 64}, + {"name": "ssh_port", "title": "SSH端口", "type": "int", "nullable": "no", "default": "22"}, + {"name": "ssh_user", "title": "SSH用户", "type": "str", "length": 64, "default": "root"}, + {"name": "role", "title": "角色", "type": "str", "length": 16}, + {"name": "status", "title": "在线状态", "type": "str", "length": 16, "nullable": "no", "default": "offline"}, + {"name": "agent_version", "title": "Agent版本", "type": "str", "length": 32}, + {"name": "last_heartbeat", "title": "最后心跳", "type": "datetime"}, + {"name": "created_at", "title": "创建时间", "type": "datetime"}, + {"name": "updated_at", "title": "更新时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_nodes_host", "idxtype": "unique", "idxfields": ["host"]} + ], + "codes": [ + {"field": "role", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='node_role'"}, + {"field": "status", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='node_status'"} + ] +} diff --git a/models/switchover_logs.json b/models/switchover_logs.json new file mode 100644 index 0000000..3781c1c --- /dev/null +++ b/models/switchover_logs.json @@ -0,0 +1,34 @@ +{ + "summary": [ + { + "name": "switchover_logs", + "title": "切换事件日志", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "日志ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "cluster_id", "title": "所属集群", "type": "str", "length": 32, "nullable": "no"}, + {"name": "event_type", "title": "事件类型", "type": "str", "length": 32, "nullable": "no"}, + {"name": "from_node_id", "title": "切换前节点", "type": "str", "length": 32}, + {"name": "to_node_id", "title": "切换后节点", "type": "str", "length": 32}, + {"name": "trigger", "title": "触发方式", "type": "str", "length": 16, "nullable": "no", "default": "auto"}, + {"name": "reason", "title": "原因说明", "type": "text"}, + {"name": "details", "title": "详细信息", "type": "text"}, + {"name": "status", "title": "执行状态", "type": "str", "length": 16, "nullable": "no", "default": "pending"}, + {"name": "started_at", "title": "开始时间", "type": "datetime"}, + {"name": "completed_at", "title": "完成时间", "type": "datetime"}, + {"name": "created_at", "title": "创建时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_sl_cluster", "idxtype": "index", "idxfields": ["cluster_id"]}, + {"name": "idx_sl_time", "idxtype": "index", "idxfields": ["created_at"]} + ], + "codes": [ + {"field": "event_type", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='switchover_event_type'"}, + {"field": "trigger", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='switchover_trigger'"}, + {"field": "status", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='switchover_status'"}, + {"field": "cluster_id", "table": "clusters", "valuefield": "id", "textfield": "name"} + ] +} diff --git a/models/sync_databases.json b/models/sync_databases.json new file mode 100644 index 0000000..360ce34 --- /dev/null +++ b/models/sync_databases.json @@ -0,0 +1,35 @@ +{ + "summary": [ + { + "name": "sync_databases", + "title": "数据库同步配置", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "配置ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "cluster_id", "title": "所属集群", "type": "str", "length": 32, "nullable": "no"}, + {"name": "db_name", "title": "数据库名", "type": "str", "length": 128, "nullable": "no"}, + {"name": "master_host", "title": "主库主机", "type": "str", "length": 255}, + {"name": "master_port", "title": "主库端口", "type": "int", "default": "3306"}, + {"name": "slave_host", "title": "从库主机", "type": "str", "length": 255}, + {"name": "slave_port", "title": "从库端口", "type": "int", "default": "3306"}, + {"name": "repl_user", "title": "复制用户", "type": "str", "length": 64}, + {"name": "repl_password", "title": "复制密码", "type": "str", "length": 255}, + {"name": "sync_status", "title": "同步状态", "type": "str", "length": 32, "default": "not_configured"}, + {"name": "seconds_behind", "title": "延迟秒数", "type": "int"}, + {"name": "last_error", "title": "最后错误", "type": "text"}, + {"name": "enabled", "title": "启用", "type": "str", "length": 1, "nullable": "no", "default": "1"}, + {"name": "created_at", "title": "创建时间", "type": "datetime"}, + {"name": "updated_at", "title": "更新时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_sync_db_cluster", "idxtype": "index", "idxfields": ["cluster_id"]}, + {"name": "idx_sync_db_name", "idxtype": "index", "idxfields": ["db_name"]} + ], + "codes": [ + {"field": "sync_status", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='sync_status'"}, + {"field": "cluster_id", "table": "clusters", "valuefield": "id", "textfield": "name"} + ] +} diff --git a/models/sync_directories.json b/models/sync_directories.json new file mode 100644 index 0000000..37a86f9 --- /dev/null +++ b/models/sync_directories.json @@ -0,0 +1,33 @@ +{ + "summary": [ + { + "name": "sync_directories", + "title": "文件目录同步配置", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "配置ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "cluster_id", "title": "所属集群", "type": "str", "length": 32, "nullable": "no"}, + {"name": "name", "title": "同步任务名称", "type": "str", "length": 128, "nullable": "no"}, + {"name": "source_host", "title": "源主机", "type": "str", "length": 255}, + {"name": "source_path", "title": "源目录路径", "type": "str", "length": 512, "nullable": "no"}, + {"name": "dest_host", "title": "目标主机", "type": "str", "length": 255}, + {"name": "dest_path", "title": "目标目录路径", "type": "str", "length": 512, "nullable": "no"}, + {"name": "sync_interval", "title": "同步间隔(秒)", "type": "int", "nullable": "no", "default": "60"}, + {"name": "exclude_patterns", "title": "排除模式", "type": "str", "length": 512}, + {"name": "last_sync_at", "title": "上次同步时间", "type": "datetime"}, + {"name": "last_sync_result", "title": "上次同步结果", "type": "str", "length": 32}, + {"name": "enabled", "title": "启用", "type": "str", "length": 1, "nullable": "no", "default": "1"}, + {"name": "created_at", "title": "创建时间", "type": "datetime"}, + {"name": "updated_at", "title": "更新时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_sync_dir_cluster", "idxtype": "index", "idxfields": ["cluster_id"]} + ], + "codes": [ + {"field": "last_sync_result", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='sync_result'"}, + {"field": "cluster_id", "table": "clusters", "valuefield": "id", "textfield": "name"} + ] +} diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..5462926 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,17 @@ +[build-system] +requires = ["setuptools>=45", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "foms" +version = "1.0.0" +description = "FOMS — Failover Management System 主备机切换管理系统 (MySQL主从复制 + 文件同步 + 故障切换)" +requires-python = ">=3.8" +dependencies = [ + "sqlor", + "bricks_for_python", +] + +[tool.setuptools.packages.find] +where = ["."] +include = ["foms*"] diff --git a/scripts/foms-agent.py b/scripts/foms-agent.py new file mode 100644 index 0000000..29ee6c7 --- /dev/null +++ b/scripts/foms-agent.py @@ -0,0 +1,260 @@ +#!/usr/bin/env python3 +""" +FOMS Agent — 部署在主机(A)和备机(B)上的守护进程 + +职责: + 1. 启动时从管理中心(C)拉取配置 + 2. 心跳上报 + 3. MySQL主从复制状态监控 + 4. 文件rsync同步 + 5. 响应切换/切回指令 +""" +import os +import sys +import json +import time +import logging +import subprocess +import urllib.request +import urllib.error +from datetime import datetime + +# ═══ Configuration ═══ +MANAGEMENT_URL = os.environ.get('FOMS_MANAGEMENT_URL', 'http://localhost:9080') +NODE_ID = os.environ.get('FOMS_NODE_ID', '') +AGENT_TOKEN = os.environ.get('FOMS_AGENT_TOKEN', '') +HEARTBEAT_INTERVAL = int(os.environ.get('FOMS_HEARTBEAT_INTERVAL', '10')) + +# ═══ Logging ═══ +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s [%(levelname)s] %(message)s', + handlers=[ + logging.FileHandler('/var/log/foms-agent.log'), + logging.StreamHandler(sys.stdout), + ] +) +log = logging.getLogger('foms-agent') + + +class FomsAgent: + def __init__(self): + self.node_id = NODE_ID + self.version = '1.0.0' + self.db_configs = [] # sync_databases configs + self.dir_configs = [] # sync_directories configs + self.node_role = '' # master or standby + + # ═══ HTTP helpers ═══ + def _api_get(self, path): + url = f"{MANAGEMENT_URL}/api/{path}" + req = urllib.request.Request(url) + req.add_header('Authorization', f'Bearer {AGENT_TOKEN}') + try: + with urllib.request.urlopen(req, timeout=10) as resp: + return json.loads(resp.read().decode()) + except Exception as e: + log.error(f"GET {url}: {e}") + return None + + def _api_post(self, path, data): + url = f"{MANAGEMENT_URL}/api/{path}" + payload = json.dumps(data).encode() + req = urllib.request.Request(url, data=payload, headers={'Content-Type': 'application/json'}) + req.add_header('Authorization', f'Bearer {AGENT_TOKEN}') + try: + with urllib.request.urlopen(req, timeout=10) as resp: + return json.loads(resp.read().decode()) + except Exception as e: + log.error(f"POST {url}: {e}") + return None + + # ═══ Heartbeat ═══ + def heartbeat(self): + data = {'node_id': self.node_id, 'version': self.version} + result = self._api_post('agent_heartbeat.dspy', data) + if result and result.get('status') == 'ok': + log.debug(f"Heartbeat OK: {result.get('server_time')}") + else: + log.warning(f"Heartbeat failed: {result}") + + # ═══ MySQL Replication Status ═══ + def check_mysql_replication(self, db_config): + """Check SHOW SLAVE STATUS for a given DB config""" + try: + cmd = [ + 'mysql', '-h', db_config.get('slave_host', '127.0.0.1'), + '-P', str(db_config.get('slave_port', 3306)), + '-u', db_config.get('repl_user', 'root'), + f"-p{db_config.get('repl_password', '')}", + '-e', 'SHOW SLAVE STATUS\\G' + ] + result = subprocess.run(cmd, capture_output=True, text=True, timeout=10) + if result.returncode != 0: + return {'id': db_config['id'], 'status': 'error', 'seconds_behind': None, 'error': result.stderr.strip()} + + # Parse SHOW SLAVE STATUS output + output = result.stdout + io_running = 'Slave_IO_Running: Yes' in output + sql_running = 'Slave_SQL_Running: Yes' in output + seconds_behind = None + for line in output.split('\n'): + if 'Seconds_Behind_Master:' in line: + val = line.split(':')[1].strip() + seconds_behind = int(val) if val.isdigit() else None + + if io_running and sql_running: + status = 'running' + error = None + elif io_running: + status = 'error' + error = 'Slave_SQL_Running: No' + elif sql_running: + status = 'error' + error = 'Slave_IO_Running: No' + else: + status = 'error' + error = 'Replication not running' + + return {'id': db_config['id'], 'status': status, 'seconds_behind': seconds_behind, 'error': error} + except Exception as e: + return {'id': db_config['id'], 'status': 'error', 'seconds_behind': None, 'error': str(e)} + + # ═══ Setup MySQL Replication ═══ + def setup_mysql_replication(self, db_config): + """On master: create replication user. On slave: CHANGE MASTER TO + START SLAVE""" + role = self.node_role + log.info(f"Setting up MySQL replication for {db_config['db_name']} (role={role})") + + if role == 'master': + # Create replication user on master + repl_user = db_config.get('repl_user', 'repl') + repl_pass = db_config.get('repl_password', '') + cmds = [ + f"CREATE USER IF NOT EXISTS '{repl_user}'@'%' IDENTIFIED BY '{repl_pass}'", + f"GRANT REPLICATION SLAVE ON *.* TO '{repl_user}'@'%'", + "FLUSH PRIVILEGES", + ] + for sql in cmds: + cmd = ['mysql', '-h', db_config.get('master_host', '127.0.0.1'), + '-P', str(db_config.get('master_port', 3306)), + '-u', 'root', f"-p{os.environ.get('MYSQL_ROOT_PASSWORD', '')}", '-e', sql] + subprocess.run(cmd, capture_output=True, timeout=10) + log.info(f" Replication user created on master") + + elif role == 'standby': + # Configure slave + # First get master status + master_host = db_config.get('master_host', '') + master_port = db_config.get('master_port', 3306) + repl_user = db_config.get('repl_user', 'repl') + repl_pass = db_config.get('repl_password', '') + + # CHANGE MASTER and START SLAVE + sql = f""" + CHANGE MASTER TO + MASTER_HOST='{master_host}', + MASTER_PORT={master_port}, + MASTER_USER='{repl_user}', + MASTER_PASSWORD='{repl_pass}', + MASTER_AUTO_POSITION=1; + START SLAVE; + """ + cmd = ['mysql', '-h', db_config.get('slave_host', '127.0.0.1'), + '-P', str(db_config.get('slave_port', 3306)), + '-u', 'root', f"-p{os.environ.get('MYSQL_ROOT_PASSWORD', '')}", '-e', sql] + result = subprocess.run(cmd, capture_output=True, text=True, timeout=15) + if result.returncode != 0: + log.error(f" Slave setup failed: {result.stderr}") + else: + log.info(f" Slave configured for {db_config['db_name']}") + + # ═══ File Sync (rsync) ═══ + def sync_directory(self, dir_config): + """Execute rsync for a directory config""" + try: + src = f"{dir_config.get('source_host', '')}:{dir_config['source_path']}" + dst = dir_config['dest_path'] + exclude = dir_config.get('exclude_patterns', '') + cmd = ['rsync', '-avz', '--delete'] + if exclude: + for pattern in exclude.split(','): + cmd.extend(['--exclude', pattern.strip()]) + cmd.extend(['-e', 'ssh', src, dst]) + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=300) + success = result.returncode == 0 + log.info(f"rsync {dir_config['name']}: {'OK' if success else 'FAILED'}") + if not success: + log.error(f" stderr: {result.stderr[:200]}") + return {'id': dir_config['id'], 'result': 'success' if success else 'failed'} + except Exception as e: + log.error(f"rsync {dir_config.get('name', '')} error: {e}") + return {'id': dir_config['id'], 'result': 'failed'} + + # ═══ Report Health ═══ + def report_health(self): + """Collect and report all health statuses""" + db_statuses = [] + for db_config in self.db_configs: + if db_config.get('enabled', '1') == '0': + continue + status = self.check_mysql_replication(db_config) + db_statuses.append(status) + + dir_statuses = [] + for dir_config in self.dir_configs: + if dir_config.get('enabled', '1') == '0': + continue + status = self.sync_directory(dir_config) + dir_statuses.append(status) + + data = { + 'node_id': self.node_id, + 'db_statuses': db_statuses, + 'dir_statuses': dir_statuses, + } + self._api_post('agent_report_health.dspy', data) + + # ═══ Fetch Config from Management ═══ + def fetch_config(self): + """Fetch sync configs from management server""" + # In production, the management would have dedicated APIs for this + # For now, the config is set via environment or manual setup + log.info(f"Agent running on node {self.node_id}") + log.info(f"Management URL: {MANAGEMENT_URL}") + log.info(f"Heartbeat interval: {HEARTBEAT_INTERVAL}s") + + # ═══ Main Loop ═══ + def run(self): + log.info(f"FOMS Agent v{self.version} starting...") + self.fetch_config() + + last_health_report = 0 + HEALTH_REPORT_INTERVAL = 60 # report health every 60s + + while True: + try: + self.heartbeat() + + # Report health (includes MySQL status + rsync) + now = time.time() + if now - last_health_report > HEALTH_REPORT_INTERVAL: + self.report_health() + last_health_report = now + + time.sleep(HEARTBEAT_INTERVAL) + except KeyboardInterrupt: + log.info("Agent stopped") + break + except Exception as e: + log.error(f"Loop error: {e}") + time.sleep(HEARTBEAT_INTERVAL) + + +if __name__ == '__main__': + if not NODE_ID: + log.error("FOMS_NODE_ID environment variable not set") + sys.exit(1) + agent = FomsAgent() + agent.run() diff --git a/scripts/load_path.py b/scripts/load_path.py new file mode 100644 index 0000000..1d3ef7d --- /dev/null +++ b/scripts/load_path.py @@ -0,0 +1,84 @@ +#!/usr/bin/env python3 +"""FOMS RBAC permission registration""" +import os, sys + +# Find Sage root for set_role_perm.py +sage_root = None +for c in [os.path.expanduser("~/repos/sage"), os.path.expanduser("~/test/sage")]: + if os.path.isdir(os.path.join(c, "py3", "bin")): + sage_root = c + break +if not sage_root: + print("ERROR: Sage root not found") + sys.exit(1) + +sys.path.insert(0, os.path.join(sage_root, "py3", "bin")) + +from set_role_perm import set_path_perm + +# Paths accessible without login +PATHS_ANY = [ + "/", + "/index.ui", + "/api/login.dspy", +] + +# Paths requiring authentication +PATHS_LOGINED = [ + "/api/clusters_create.dspy", + "/api/clusters_update.dspy", + "/api/clusters_delete.dspy", + "/api/nodes_create.dspy", + "/api/nodes_update.dspy", + "/api/nodes_delete.dspy", + "/api/sync_databases_create.dspy", + "/api/sync_databases_update.dspy", + "/api/sync_databases_delete.dspy", + "/api/sync_directories_create.dspy", + "/api/sync_directories_update.dspy", + "/api/sync_directories_delete.dspy", + "/api/agent_heartbeat.dspy", + "/api/agent_report_health.dspy", + "/api/trigger_switchover.dspy", + "/api/trigger_switchback.dspy", + "/api/dashboard_stats.dspy", + # CRUD list pages + "/clusters_list", + "/clusters_list/index.ui", + "/nodes_list", + "/nodes_list/index.ui", + "/sync_databases_list", + "/sync_databases_list/index.ui", + "/sync_directories_list", + "/sync_directories_list/index.ui", + # CRUD generated files + "/clusters_list/get_clusters_list.dspy", + "/clusters_list/add_clusters_list.dspy", + "/clusters_list/update_clusters_list.dspy", + "/clusters_list/delete_clusters_list.dspy", + "/nodes_list/get_nodes_list.dspy", + "/nodes_list/add_nodes_list.dspy", + "/nodes_list/update_nodes_list.dspy", + "/nodes_list/delete_nodes_list.dspy", + "/sync_databases_list/get_sync_databases_list.dspy", + "/sync_databases_list/add_sync_databases_list.dspy", + "/sync_databases_list/update_sync_databases_list.dspy", + "/sync_databases_list/delete_sync_databases_list.dspy", + "/sync_directories_list/get_sync_directories_list.dspy", + "/sync_directories_list/add_sync_directories_list.dspy", + "/sync_directories_list/update_sync_directories_list.dspy", + "/sync_directories_list/delete_sync_directories_list.dspy", +] + +# Agent heartbeat endpoint: no auth (agents use token) +PATHS_ANY.append("/api/agent_heartbeat.dspy") +PATHS_ANY.append("/api/agent_report_health.dspy") + +if __name__ == '__main__': + for p in PATHS_ANY: + set_path_perm("any", p) + print(f" any {p}") + for p in PATHS_LOGINED: + set_path_perm("logined", p) + print(f" logined {p}") + print("Done. Restart FOMS to apply.") diff --git a/start.sh b/start.sh new file mode 100755 index 0000000..a00aee6 --- /dev/null +++ b/start.sh @@ -0,0 +1,6 @@ +#!/usr/bin/env bash +cd "$(dirname "$0")" +[ -f .env ] && source .env +nohup $PWD/py3/bin/python $PWD/app/foms.py -p ${FOMS_PORT:-9080} -w $PWD > $PWD/logs/foms.out 2>&1 & +echo $! > foms.pid +echo "FOMS started (PID: $(cat foms.pid)) on port ${FOMS_PORT:-9080}" diff --git a/stop.sh b/stop.sh new file mode 100755 index 0000000..192f6fa --- /dev/null +++ b/stop.sh @@ -0,0 +1,9 @@ +#!/usr/bin/env bash +cd "$(dirname "$0")" +if [ -f foms.pid ]; then + pid=$(cat foms.pid) + kill $pid 2>/dev/null && echo "FOMS stopped (PID: $pid)" || echo "FOMS not running" + rm -f foms.pid +else + echo "No PID file found" +fi diff --git a/wwwroot/api/agent_heartbeat.dspy b/wwwroot/api/agent_heartbeat.dspy new file mode 100644 index 0000000..9adc9aa --- /dev/null +++ b/wwwroot/api/agent_heartbeat.dspy @@ -0,0 +1,7 @@ +# Agent heartbeat endpoint +node_id = params_kw.get('node_id', '') +version = params_kw.get('version', '') +if not node_id: + return json.dumps({'status': 'error', 'message': 'node_id required'}, ensure_ascii=False) +result = await agent_heartbeat(request, params_kw) +return json.dumps(result, ensure_ascii=False) diff --git a/wwwroot/api/agent_report_health.dspy b/wwwroot/api/agent_report_health.dspy new file mode 100644 index 0000000..b8e1561 --- /dev/null +++ b/wwwroot/api/agent_report_health.dspy @@ -0,0 +1,3 @@ +# Agent health report endpoint (sync statuses) +result = await agent_report_health(request, params_kw) +return json.dumps(result, ensure_ascii=False) diff --git a/wwwroot/api/clusters_create.dspy b/wwwroot/api/clusters_create.dspy new file mode 100644 index 0000000..829ef32 --- /dev/null +++ b/wwwroot/api/clusters_create.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'create_cluster', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "create_cluster not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/clusters_delete.dspy b/wwwroot/api/clusters_delete.dspy new file mode 100644 index 0000000..1f3c94d --- /dev/null +++ b/wwwroot/api/clusters_delete.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'delete_clusters', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "delete_clusters not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/clusters_update.dspy b/wwwroot/api/clusters_update.dspy new file mode 100644 index 0000000..dc554f0 --- /dev/null +++ b/wwwroot/api/clusters_update.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'update_clusters', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "update_clusters not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/dashboard_stats.dspy b/wwwroot/api/dashboard_stats.dspy new file mode 100644 index 0000000..a5ae91e --- /dev/null +++ b/wwwroot/api/dashboard_stats.dspy @@ -0,0 +1,3 @@ +# Dashboard stats API +stats = await get_dashboard_stats(request) +return json.dumps(stats, ensure_ascii=False, default=str) diff --git a/wwwroot/api/nodes_create.dspy b/wwwroot/api/nodes_create.dspy new file mode 100644 index 0000000..4f1dafa --- /dev/null +++ b/wwwroot/api/nodes_create.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'create_nodes', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "create_nodes not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/nodes_delete.dspy b/wwwroot/api/nodes_delete.dspy new file mode 100644 index 0000000..8349e4d --- /dev/null +++ b/wwwroot/api/nodes_delete.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'delete_nodes', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "delete_nodes not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/nodes_update.dspy b/wwwroot/api/nodes_update.dspy new file mode 100644 index 0000000..f2640f7 --- /dev/null +++ b/wwwroot/api/nodes_update.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'update_nodes', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "update_nodes not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_databases_create.dspy b/wwwroot/api/sync_databases_create.dspy new file mode 100644 index 0000000..7056842 --- /dev/null +++ b/wwwroot/api/sync_databases_create.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'create_sync_db', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "create_sync_db not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_databases_delete.dspy b/wwwroot/api/sync_databases_delete.dspy new file mode 100644 index 0000000..0781ed9 --- /dev/null +++ b/wwwroot/api/sync_databases_delete.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'delete_sync_db', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "delete_sync_db not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_databases_update.dspy b/wwwroot/api/sync_databases_update.dspy new file mode 100644 index 0000000..841c448 --- /dev/null +++ b/wwwroot/api/sync_databases_update.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'update_sync_db', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "update_sync_db not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_directories_create.dspy b/wwwroot/api/sync_directories_create.dspy new file mode 100644 index 0000000..1d276e8 --- /dev/null +++ b/wwwroot/api/sync_directories_create.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'create_sync_dir', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "create_sync_dir not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_directories_delete.dspy b/wwwroot/api/sync_directories_delete.dspy new file mode 100644 index 0000000..cab869b --- /dev/null +++ b/wwwroot/api/sync_directories_delete.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'delete_sync_dir', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "delete_sync_dir not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/sync_directories_update.dspy b/wwwroot/api/sync_directories_update.dspy new file mode 100644 index 0000000..8ccd7d1 --- /dev/null +++ b/wwwroot/api/sync_directories_update.dspy @@ -0,0 +1,10 @@ +# Cluster create wrapper +import json +from ahserver.serverenv import ServerEnv +env = ServerEnv() +create_func = getattr(env, 'update_sync_dir', None) +if create_func is None: + print(json.dumps({"widgettype": "Message", "options": {"title": "Error", "message": "update_sync_dir not found", "type": "error"}})) +else: + result = await create_func(request, params_kw) + print(result) diff --git a/wwwroot/api/trigger_switchback.dspy b/wwwroot/api/trigger_switchback.dspy new file mode 100644 index 0000000..bb9216f --- /dev/null +++ b/wwwroot/api/trigger_switchback.dspy @@ -0,0 +1,6 @@ +# Trigger switchback +cluster_id = params_kw.get('cluster_id', '') +if not cluster_id: + return json.dumps({'status': 'error', 'message': 'cluster_id required'}, ensure_ascii=False) +result = await trigger_switchback(request, params_kw) +return json.dumps(result, ensure_ascii=False) diff --git a/wwwroot/api/trigger_switchover.dspy b/wwwroot/api/trigger_switchover.dspy new file mode 100644 index 0000000..570dab8 --- /dev/null +++ b/wwwroot/api/trigger_switchover.dspy @@ -0,0 +1,6 @@ +# Trigger switchover +cluster_id = params_kw.get('cluster_id', '') +if not cluster_id: + return json.dumps({'status': 'error', 'message': 'cluster_id required'}, ensure_ascii=False) +result = await trigger_switchover(request, params_kw) +return json.dumps(result, ensure_ascii=False) diff --git a/wwwroot/index.ui b/wwwroot/index.ui new file mode 100644 index 0000000..6ba83b4 --- /dev/null +++ b/wwwroot/index.ui @@ -0,0 +1,87 @@ +{ + "widgettype": "VBox", + "options": {"width": "100%", "height": "100%", "padding": "24px", "backgroundColor": "#F5F7FA"}, + "subwidgets": [ + { + "widgettype": "Text", + "options": {"text": "FOMS — 主备切换管理系统", "fontSize": "28px", "fontWeight": "bold", "marginBottom": "8px"} + }, + { + "widgettype": "Text", + "options": {"text": "Failover Management System · 实时监控 · 自动切换 · 数据同步", "fontSize": "14px", "color": "#666", "marginBottom": "24px"} + }, + { + "widgettype": "ResponsableBox", + "options": {"gap": "16px", "minWidth": "220px", "marginBottom": "24px"}, + "subwidgets": [ + { + "widgettype": "VBox", + "options": {"backgroundColor": "#FFFFFF", "padding": "20px", "borderRadius": "8px", "cursor": "pointer"}, + "binds": [{"wid": "self", "event": "click", "actiontype": "urlwidget", "target": "app.foms_content", "options": {"url": "{{entire_url('clusters_list')}}"}, "mode": "replace"}], + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "集群管理", "fontSize": "16px", "fontWeight": "bold", "marginBottom": "8px"}}, + {"widgettype": "RefreshWidget", "options": {"interval": 30000, "data_url": "{{entire_url('api/dashboard_stats.dspy')}}"}, + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "{{data.total_clusters}} 个集群", "fontSize": "14px", "color": "#666"}} + ] + } + ] + }, + { + "widgettype": "VBox", + "options": {"backgroundColor": "#FFFFFF", "padding": "20px", "borderRadius": "8px", "cursor": "pointer"}, + "binds": [{"wid": "self", "event": "click", "actiontype": "urlwidget", "target": "app.foms_content", "options": {"url": "{{entire_url('nodes_list')}}"}, "mode": "replace"}], + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "节点管理", "fontSize": "16px", "fontWeight": "bold", "marginBottom": "8px"}}, + {"widgettype": "RefreshWidget", "options": {"interval": 30000, "data_url": "{{entire_url('api/dashboard_stats.dspy')}}"}, + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "{{data.nodes_online}} 在线 / {{data.nodes_offline}} 离线", "fontSize": "14px", "color": "#666"}} + ] + } + ] + }, + { + "widgettype": "VBox", + "options": {"backgroundColor": "#FFFFFF", "padding": "20px", "borderRadius": "8px", "cursor": "pointer"}, + "binds": [{"wid": "self", "event": "click", "actiontype": "urlwidget", "target": "app.foms_content", "options": {"url": "{{entire_url('sync_databases_list')}}"}, "mode": "replace"}], + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "数据库同步", "fontSize": "16px", "fontWeight": "bold", "marginBottom": "8px"}}, + {"widgettype": "RefreshWidget", "options": {"interval": 30000, "data_url": "{{entire_url('api/dashboard_stats.dspy')}}"}, + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "{{data.dbs_synced}} 正常 / {{data.dbs_error}} 异常", "fontSize": "14px", "color": "#666"}} + ] + } + ] + }, + { + "widgettype": "VBox", + "options": {"backgroundColor": "#FFFFFF", "padding": "20px", "borderRadius": "8px", "cursor": "pointer"}, + "binds": [{"wid": "self", "event": "click", "actiontype": "urlwidget", "target": "app.foms_content", "options": {"url": "{{entire_url('sync_directories_list')}}"}, "mode": "replace"}], + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "文件同步", "fontSize": "16px", "fontWeight": "bold", "marginBottom": "8px"}}, + {"widgettype": "RefreshWidget", "options": {"interval": 30000, "data_url": "{{entire_url('api/dashboard_stats.dspy')}}"}, + "subwidgets": [ + {"widgettype": "Text", "options": {"text": "{{data.dirs_synced}} 同步正常", "fontSize": "14px", "color": "#666"}} + ] + } + ] + } + ] + }, + { + "widgettype": "VBox", + "id": "foms_content", + "options": {"width": "100%", "flex": "1", "backgroundColor": "#FFFFFF", "borderRadius": "8px", "padding": "20px", "minHeight": "400px"}, + "subwidgets": [ + { + "widgettype": "Text", + "options": {"text": "欢迎使用 FOMS 主备切换管理系统", "fontSize": "18px", "color": "#999", "textAlign": "center", "marginTop": "80px"} + }, + { + "widgettype": "Text", + "options": {"text": "点击上方卡片进入各功能模块", "fontSize": "14px", "color": "#BBB", "textAlign": "center"} + } + ] + } + ] +}