Compare commits

...

No commits in common. "master" and "main" have entirely different histories.
master ... main

53 changed files with 1 additions and 2214 deletions

20
.gitignore vendored
View File

@ -1,20 +0,0 @@
__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/

View File

@ -1,61 +1,2 @@
# FOMS — Failover Management System # foms
主备机切换管理系统。运行在管理中心 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
```

View File

@ -1,32 +0,0 @@
#!/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)

View File

@ -1,12 +0,0 @@
"""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

View File

@ -1,51 +0,0 @@
#!/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 ==="

View File

@ -1,42 +0,0 @@
{
"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"]
]
}
}

View File

@ -1,12 +0,0 @@
from .init import (
create_cluster, update_cluster, delete_cluster,
create_node, update_node, delete_node,
create_sync_db, update_sync_db, delete_sync_db,
create_sync_dir, update_sync_dir, delete_sync_dir,
agent_heartbeat, agent_report_health,
trigger_switchover, trigger_switchback,
get_dashboard_stats,
agent_report_metrics, agent_push_logs,
get_node_detail, get_node_metrics_history,
exec_remote_command, agent_poll_commands, agent_report_command_result,
)

View File

@ -1,560 +0,0 @@
"""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 [])],
}
# ═══════════════════════════════════════════
# Metrics & Log Collection
# ═══════════════════════════════════════════
async def agent_report_metrics(request, params_kw):
"""Agent上报系统指标CPU/内存/磁盘/负载/网络)"""
env = ServerEnv()
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
metrics = params_kw.get('metrics', {}) if hasattr(params_kw, 'get') else {}
if isinstance(metrics, str):
metrics = json.loads(metrics)
async with get_sor_context(env, DBNAME) as sor:
await sor.C('node_metrics', {
'id': getID(),
'node_id': node_id,
'cpu_percent': metrics.get('cpu_percent'),
'mem_percent': metrics.get('mem_percent'),
'mem_used_mb': metrics.get('mem_used_mb'),
'mem_total_mb': metrics.get('mem_total_mb'),
'disk_percent': metrics.get('disk_percent'),
'disk_used_gb': metrics.get('disk_used_gb'),
'disk_path': metrics.get('disk_path', '/'),
'load_1m': metrics.get('load_1m'),
'load_5m': metrics.get('load_5m'),
'load_15m': metrics.get('load_15m'),
'net_rx_mb': metrics.get('net_rx_mb'),
'net_tx_mb': metrics.get('net_tx_mb'),
'process_count': metrics.get('process_count'),
'collected_at': curDateString(),
})
return {'status': 'ok'}
async def agent_push_logs(request, params_kw):
"""Agent推送日志条目到管理中心"""
env = ServerEnv()
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
logs = params_kw.get('logs', []) if hasattr(params_kw, 'get') else []
if isinstance(logs, str):
logs = json.loads(logs)
if not logs:
return {'status': 'ok', 'count': 0}
async with get_sor_context(env, DBNAME) as sor:
count = 0
for entry in logs[:100]: # max 100 per batch
await sor.C('node_logs', {
'id': getID(),
'node_id': node_id,
'log_source': entry.get('source', 'syslog'),
'log_level': entry.get('level', 'INFO'),
'message': entry.get('message', '')[:2000],
'logged_at': entry.get('time', curDateString()),
})
count += 1
return {'status': 'ok', 'count': count}
async def get_node_detail(request, params_kw):
"""获取节点详情(实时指标 + 最新日志 + 基本信息)"""
env = ServerEnv()
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
async with get_sor_context(env, DBNAME) as sor:
node = await sor.R('nodes', {'id': node_id})
metrics = await sor.sqlExe(
'SELECT * FROM node_metrics WHERE node_id=${nid}$ ORDER BY collected_at DESC LIMIT 1',
{'nid': node_id}
)
logs = await sor.sqlExe(
'SELECT log_source, log_level, message, logged_at FROM node_logs WHERE node_id=${nid}$ ORDER BY logged_at DESC LIMIT 50',
{'nid': node_id}
)
node_info = {}
if node:
n = node[0]
node_info = {
'id': getattr(n, 'id', ''),
'host': getattr(n, 'host', ''),
'name': getattr(n, 'name', ''),
'ip': getattr(n, 'ip', ''),
'role': getattr(n, 'role', ''),
'status': getattr(n, 'status', ''),
'agent_version': getattr(n, 'agent_version', ''),
'last_heartbeat': str(getattr(n, 'last_heartbeat', '')),
}
return {
'node': node_info,
'latest_metrics': dict(metrics[0]) if metrics else {},
'recent_logs': [{
'source': getattr(l, 'log_source', ''),
'level': getattr(l, 'log_level', ''),
'message': getattr(l, 'message', ''),
'time': str(getattr(l, 'logged_at', '')),
} for l in (logs or [])],
}
async def get_node_metrics_history(request, params_kw):
"""获取节点历史指标(用于图表)"""
env = ServerEnv()
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
hours = int(params_kw.get('hours', 1)) if hasattr(params_kw, 'get') else 1
async with get_sor_context(env, DBNAME) as sor:
rows = await sor.sqlExe(
'SELECT cpu_percent, mem_percent, disk_percent, load_1m, collected_at '
'FROM node_metrics WHERE node_id=${nid}$ AND collected_at > DATE_SUB(NOW(), INTERVAL ${h}$ HOUR) '
'ORDER BY collected_at ASC LIMIT 200',
{'nid': node_id, 'h': hours}
)
return {'metrics': [{
'cpu': getattr(r, 'cpu_percent', 0),
'mem': getattr(r, 'mem_percent', 0),
'disk': getattr(r, 'disk_percent', 0),
'load': getattr(r, 'load_1m', 0),
'time': str(getattr(r, 'collected_at', '')),
} for r in (rows or [])]}
# ═══════════════════════════════════════════
# Remote Command Execution
# ═══════════════════════════════════════════
async def exec_remote_command(request, params_kw):
"""向Agent下发远程命令Agent执行后返回结果"""
# Commands are pushed to agent via its next heartbeat poll, or via a pending queue
# For simplicity, this stores the command; agent reads it via poll_commands
env = ServerEnv()
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
command = params_kw.get('command', '') if hasattr(params_kw, 'get') else ''
if not node_id or not command:
return {'status': 'error', 'message': 'node_id and command required'}
# Store pending command (simple approach: use sync_databases or a new table)
# For MVP: just validate and return that command was queued
# Agent checks for commands via agent_poll_commands below
cmd_id = getID()
now = curDateString()
async with get_sor_context(env, DBNAME) as sor:
await sor.C('node_logs', {
'id': cmd_id,
'node_id': node_id,
'log_source': 'foms_cmd',
'log_level': 'INFO',
'message': f'[PENDING] {command}',
'logged_at': now,
})
return {
'status': 'ok',
'message': '命令已下发等待Agent执行',
'cmd_id': cmd_id,
}
async def agent_poll_commands(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:
rows = await sor.sqlExe(
"SELECT id, message FROM node_logs WHERE node_id=${nid}$ AND log_source='foms_cmd' AND message LIKE '[PENDING]%' ORDER BY logged_at ASC LIMIT 5",
{'nid': node_id}
)
commands = []
for r in (rows or []):
msg = getattr(r, 'message', '')
cmd = msg.replace('[PENDING] ', '') if '[PENDING]' in msg else msg
commands.append({'id': getattr(r, 'id', ''), 'command': cmd})
return {'commands': commands}
async def agent_report_command_result(request, params_kw):
"""Agent上报命令执行结果"""
env = ServerEnv()
cmd_id = params_kw.get('cmd_id', '') if hasattr(params_kw, 'get') else ''
node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else ''
output = params_kw.get('output', '') if hasattr(params_kw, 'get') else ''
exit_code = params_kw.get('exit_code', -1) if hasattr(params_kw, 'get') else -1
async with get_sor_context(env, DBNAME) as sor:
await sor.U('node_logs', {
'id': cmd_id,
'log_source': 'foms_cmd',
'log_level': 'INFO' if exit_code == 0 else 'ERROR',
'message': f'[EXIT:{exit_code}] {output[:2000]}',
})
return {'status': 'ok'}
# ═══════════════════════════════════════════
# 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
# Metrics & Logs
env.agent_report_metrics = agent_report_metrics
env.agent_push_logs = agent_push_logs
env.get_node_detail = get_node_detail
env.get_node_metrics_history = get_node_metrics_history
# Remote Command
env.exec_remote_command = exec_remote_command
env.agent_poll_commands = agent_poll_commands
env.agent_report_command_result = agent_report_command_result
return True

View File

@ -1,171 +0,0 @@
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
- id: org_type
name: 机构类型
hierarchy_flg: 0
- id: user_status
name: 用户状态
hierarchy_flg: 0
- id: log_level
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: 失败
- id: org_type_foms
parentid: org_type
k: foms_org
v: FOMS运维组织
- id: user_status_active
parentid: user_status
k: 0
v: 正常
- id: user_status_disabled
parentid: user_status
k: 1
v: 停用
- id: log_level_info
parentid: log_level
k: INFO
v: 信息
- id: log_level_warn
parentid: log_level
k: WARN
v: 警告
- id: log_level_error
parentid: log_level
k: ERROR
v: 错误
- id: log_level_crit
parentid: log_level
k: CRIT
v: 严重

View File

@ -1,22 +0,0 @@
{
"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')}}"
}
}
}

View File

@ -1,21 +0,0 @@
{
"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')}}"
}
}
}

View File

@ -1,22 +0,0 @@
{
"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')}}"
}
}
}

View File

@ -1,22 +0,0 @@
{
"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')}}"
}
}
}

View File

@ -1,30 +0,0 @@
{
"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"}
]
}

View File

@ -1,26 +0,0 @@
{
"summary": [
{
"name": "node_logs",
"title": "节点日志收集",
"primary": ["id"],
"catelog": "entity"
}
],
"fields": [
{"name": "id", "title": "日志ID", "type": "str", "length": 32, "nullable": "no"},
{"name": "node_id", "title": "节点ID", "type": "str", "length": 32, "nullable": "no"},
{"name": "log_source", "title": "日志来源", "type": "str", "length": 64},
{"name": "log_level", "title": "日志级别", "type": "str", "length": 16},
{"name": "message", "title": "日志内容", "type": "text"},
{"name": "logged_at", "title": "日志时间", "type": "datetime", "nullable": "no"}
],
"indexes": [
{"name": "idx_nl_node", "idxtype": "index", "idxfields": ["node_id"]},
{"name": "idx_nl_time", "idxtype": "index", "idxfields": ["logged_at"]},
{"name": "idx_nl_level", "idxtype": "index", "idxfields": ["log_level"]}
],
"codes": [
{"field": "log_level", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='log_level'"}
]
}

View File

@ -1,32 +0,0 @@
{
"summary": [
{
"name": "node_metrics",
"title": "节点系统指标",
"primary": ["id"],
"catelog": "indication"
}
],
"fields": [
{"name": "id", "title": "记录ID", "type": "str", "length": 32, "nullable": "no"},
{"name": "node_id", "title": "节点ID", "type": "str", "length": 32, "nullable": "no"},
{"name": "cpu_percent", "title": "CPU使用率(%)", "type": "float", "length": 5, "dec": 1},
{"name": "mem_percent", "title": "内存使用率(%)", "type": "float", "length": 5, "dec": 1},
{"name": "mem_used_mb", "title": "已用内存(MB)", "type": "int"},
{"name": "mem_total_mb", "title": "总内存(MB)", "type": "int"},
{"name": "disk_percent", "title": "磁盘使用率(%)", "type": "float", "length": 5, "dec": 1},
{"name": "disk_used_gb", "title": "已用磁盘(GB)", "type": "float", "length": 8, "dec": 1},
{"name": "disk_path", "title": "磁盘挂载点", "type": "str", "length": 128},
{"name": "load_1m", "title": "1分钟负载", "type": "float", "length": 5, "dec": 2},
{"name": "load_5m", "title": "5分钟负载", "type": "float", "length": 5, "dec": 2},
{"name": "load_15m", "title": "15分钟负载", "type": "float", "length": 5, "dec": 2},
{"name": "net_rx_mb", "title": "网络接收(MB)", "type": "float", "length": 10, "dec": 1},
{"name": "net_tx_mb", "title": "网络发送(MB)", "type": "float", "length": 10, "dec": 1},
{"name": "process_count", "title": "进程数", "type": "int"},
{"name": "collected_at", "title": "采集时间", "type": "datetime", "nullable": "no"}
],
"indexes": [
{"name": "idx_nm_node", "idxtype": "index", "idxfields": ["node_id"]},
{"name": "idx_nm_time", "idxtype": "index", "idxfields": ["collected_at"]}
]
}

View File

@ -1,31 +0,0 @@
{
"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'"}
]
}

View File

@ -1,34 +0,0 @@
{
"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"}
]
}

View File

@ -1,35 +0,0 @@
{
"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"}
]
}

View File

@ -1,33 +0,0 @@
{
"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"}
]
}

View File

@ -1,17 +0,0 @@
[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*"]

View File

@ -1,47 +0,0 @@
#!/usr/bin/env bash
# FOMS Agent 一键部署 — 在 A/B 机器上执行
set -e
INSTALL_DIR=${1:-/opt/foms-agent}
MANAGEMENT_URL=${2:-http://192.168.1.100:9080}
NODE_ID=${3:-}
AGENT_TOKEN=${4:-}
if [ -z "$NODE_ID" ]; then
echo "用法: $0 <安装目录> <管理端URL> <节点ID> <Agent令牌>"
echo "示例: $0 /opt/foms-agent http://192.168.1.100:9080 node_a_001 mytoken123"
exit 1
fi
echo "=== FOMS Agent 部署 ==="
echo " 安装目录: $INSTALL_DIR"
echo " 管理端: $MANAGEMENT_URL"
echo " 节点ID: $NODE_ID"
# 1. 创建目录
mkdir -p "$INSTALL_DIR"
# 2. 复制 agent 脚本
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)"
cp "$SCRIPT_DIR/foms-agent.py" "$INSTALL_DIR/"
# 3. 创建 .env
cat > "$INSTALL_DIR/.env" << EOF
FOMS_MANAGEMENT_URL=$MANAGEMENT_URL
FOMS_NODE_ID=$NODE_ID
FOMS_AGENT_TOKEN=$AGENT_TOKEN
FOMS_HEARTBEAT_INTERVAL=10
EOF
# 4. 安装 systemd 服务
cp "$SCRIPT_DIR/foms-agent.service" /etc/systemd/system/foms-agent.service
sed -i "s|/opt/foms-agent|$INSTALL_DIR|g" /etc/systemd/system/foms-agent.service
# 5. 启用并启动
systemctl daemon-reload
systemctl enable foms-agent
systemctl restart foms-agent
echo "=== 部署完成 ==="
echo " systemctl status foms-agent # 查看状态"
echo " journalctl -u foms-agent -f # 查看日志"

View File

@ -1,360 +0,0 @@
#!/usr/bin/env python3
"""
FOMS Agent 部署在主机(A)和备机(B)上的守护进程
职责:
1. 启动时从管理中心(C)拉取配置
2. 心跳上报
3. MySQL主从复制状态监控
4. 文件rsync同步
5. 响应切换/切回指令
6. 系统指标采集 (CPU/内存/磁盘/网络)
7. 系统日志推送到管理中心
8. 远程命令执行
"""
import os
import sys
import json
import time
import logging
import subprocess
import urllib.request
import urllib.error
from datetime import datetime
try:
import psutil
HAS_PSUTIL = True
except ImportError:
HAS_PSUTIL = False
# ═══ 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...")
if HAS_PSUTIL:
log.info(" psutil: available (system metrics enabled)")
else:
log.warning(" psutil: NOT installed — pip install psutil")
self.fetch_config()
last_health_report = 0
last_metrics_report = 0
last_command_poll = 0
HEALTH_REPORT_INTERVAL = 60
METRICS_INTERVAL = 30
while True:
try:
self.heartbeat()
now = time.time()
# Report health (MySQL status + rsync)
if now - last_health_report > HEALTH_REPORT_INTERVAL:
self.report_health()
last_health_report = now
# Collect & report system metrics
if now - last_metrics_report > METRICS_INTERVAL:
self.report_metrics()
last_metrics_report = now
# Poll for remote commands
if now - last_command_poll > 15:
self.poll_and_execute_commands()
last_command_poll = 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)
# ═══ System Metrics ═══
def collect_metrics(self):
"""Collect system metrics via psutil"""
if not HAS_PSUTIL:
return {}
try:
cpu = psutil.cpu_percent(interval=0.5)
mem = psutil.virtual_memory()
disk = psutil.disk_usage('/')
load = os.getloadavg() if hasattr(os, 'getloadavg') else (0, 0, 0)
net = psutil.net_io_counters()
uptime = int(time.time() - psutil.boot_time())
return {
'cpu_percent': round(cpu, 1),
'mem_percent': round(mem.percent, 1),
'mem_used_mb': round(mem.used / 1024 / 1024),
'mem_total_mb': round(mem.total / 1024 / 1024),
'disk_percent': round(disk.percent, 1),
'disk_used_gb': round(disk.used / 1024 / 1024 / 1024, 1),
'disk_path': '/',
'load_1m': round(load[0], 2),
'load_5m': round(load[1], 2),
'load_15m': round(load[2], 2),
'net_rx_mb': round(net.bytes_recv / 1024 / 1024, 1),
'net_tx_mb': round(net.bytes_sent / 1024 / 1024, 1),
'process_count': len(psutil.pids()),
}
except Exception as e:
log.error(f"collect_metrics error: {e}")
return {}
def report_metrics(self):
"""Send metrics to management"""
metrics = self.collect_metrics()
if not metrics:
return
data = {'node_id': self.node_id, 'metrics': metrics}
result = self._api_post('agent_report_metrics.dspy', data)
if result and result.get('status') == 'ok':
log.debug(f"Metrics reported: CPU={metrics.get('cpu_percent')}% MEM={metrics.get('mem_percent')}%")
else:
log.warning(f"Metrics report failed: {result}")
# ═══ Remote Command Execution ═══
def poll_and_execute_commands(self):
"""Poll management for pending commands, execute, report results"""
result = self._api_get('agent_poll_commands.dspy?node_id=' + self.node_id)
if not result or not result.get('commands'):
return
for cmd in result['commands']:
cmd_id = cmd.get('id', '')
command = cmd.get('command', '')
log.info(f"Executing command [{cmd_id}]: {command}")
try:
proc = subprocess.run(command, shell=True, capture_output=True, text=True, timeout=30)
output = (proc.stdout + proc.stderr)[:2000]
self._api_post('agent_report_command_result.dspy', {
'node_id': self.node_id,
'cmd_id': cmd_id,
'output': output,
'exit_code': proc.returncode,
})
log.info(f"Command [{cmd_id}] done: exit={proc.returncode}")
except subprocess.TimeoutExpired:
self._api_post('agent_report_command_result.dspy', {
'node_id': self.node_id, 'cmd_id': cmd_id,
'output': 'TIMEOUT (30s)', 'exit_code': -1,
})
except Exception as e:
self._api_post('agent_report_command_result.dspy', {
'node_id': self.node_id, 'cmd_id': cmd_id,
'output': str(e), 'exit_code': -1,
})
if __name__ == '__main__':
if not NODE_ID:
log.error("FOMS_NODE_ID environment variable not set")
sys.exit(1)
agent = FomsAgent()
agent.run()

View File

@ -1,21 +0,0 @@
[Unit]
Description=FOMS Agent - Failover Management System Agent
Documentation=https://git.opencomputing.cn/yumoqing/foms
After=network.target mysqld.service
[Service]
Type=simple
User=root
EnvironmentFile=/opt/foms-agent/.env
WorkingDirectory=/opt/foms-agent
ExecStart=/usr/bin/python3 /opt/foms-agent/foms-agent.py
Restart=always
RestartSec=10
StandardOutput=append:/var/log/foms-agent.log
StandardError=append:/var/log/foms-agent.log
# 防止无限重启耗尽资源
StartLimitInterval=300
StartLimitBurst=5
[Install]
WantedBy=multi-user.target

View File

@ -1,155 +0,0 @@
#!/usr/bin/env python3
"""
FOMS RBAC 初始化 创建机构类型机构角色用户
用法: cd /d/ymq/repos/foms && source py3/bin/activate && python scripts/init_rbac.py
"""
import os, sys, asyncio
# Add project root
root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.insert(0, root)
from appPublic.uniqueID import getID
from appPublic.timeUtils import curDateString
from sqlor.dbpools import DBPools
from ahserver.serverenv import ServerEnv
DBNAME = 'foms'
ORG_ID = 'foms_org_001'
ROLES = {
'foms_superadmin': '超级管理员',
'foms_operator': '运维操作员',
'foms_viewer': '只读观察者',
}
USERS = {
'admin': {'name': '系统管理员', 'password': 'Admin@123', 'role': 'foms_superadmin'},
'operator': {'name': '运维操作员', 'password': 'Ops@123456', 'role': 'foms_operator'},
}
async def main():
env = ServerEnv()
dbname = env.get_module_dbname('rbac') if hasattr(env, 'get_module_dbname') else DBNAME
db = DBPools()
async with db.sqlorContext(dbname) as sor:
# 1. Create organization
existing = await sor.R('organization', {'id': ORG_ID})
if existing:
print(f"[SKIP] organization {ORG_ID} already exists")
else:
await sor.C('organization', {
'id': ORG_ID,
'orgname': 'FOMS运维中心',
'orgabbr': 'FOMS',
'org_type': 'foms_org',
})
print(f"[OK] organization created: FOMS运维中心")
# 2. Create roles
for role_id, role_name in ROLES.items():
existing = await sor.R('role', {'id': role_id})
if existing:
print(f"[SKIP] role {role_id} exists")
continue
await sor.C('role', {
'id': role_id,
'orgtypeid': 'foms_org',
'name': role_name,
})
print(f"[OK] role created: {role_id} ({role_name})")
# 3. Create users
for username, info in USERS.items():
uid = f'user_{username}'
existing = await sor.R('users', {'username': username})
if existing:
print(f"[SKIP] user {username} exists")
else:
encoded_pw = env.password_encode(info['password'])
await sor.C('users', {
'id': uid,
'username': username,
'name': info['name'],
'nick_name': info['name'],
'password': encoded_pw,
'orgid': ORG_ID,
'user_status': '0',
'created_at': curDateString(),
})
print(f"[OK] user created: {username}")
# 4. Assign role to user
ur_id = f'ur_{username}_{info["role"]}'
existing_ur = await sor.R('userrole', {'userid': uid, 'roleid': info['role']})
if existing_ur:
print(f"[SKIP] userrole {username}{info['role']} exists")
continue
await sor.C('userrole', {
'id': ur_id,
'userid': uid,
'roleid': info['role'],
})
print(f"[OK] userrole: {username}{info['role']}")
# 5. Assign permissions to roles
# Get all FOMS-specific permissions (paths not from rbac/appbase)
all_perms = await sor.sqlExe(
"SELECT id, path FROM permission WHERE path IN ("
"SELECT path FROM permission WHERE path LIKE '/api/%' OR path LIKE '/%_list%'"
") ORDER BY path", {}
)
# Most paths start with / for FOMS root-mapped module
foms_perms = [p for p in all_perms if hasattr(p, 'path')]
if not foms_perms:
print("[WARN] No FOMS permissions found — run scripts/load_path.py first")
else:
# superadmin: all FOMS paths
# operator: all paths (can trigger switchover)
# viewer: read-only (no create/update/delete)
for perm in foms_perms:
pid = getattr(perm, 'id', '')
path = getattr(perm, 'path', '')
is_write = any(x in path for x in ['create', 'update', 'delete', 'trigger'])
is_agent = 'agent_' in path
# superadmin gets everything
await _add_roleperm(sor, 'foms_superadmin', pid, path)
# operator gets everything except agent paths (internal)
if not is_agent:
await _add_roleperm(sor, 'foms_operator', pid, path)
# viewer gets read-only
if not is_write and not is_agent:
await _add_roleperm(sor, 'foms_viewer', pid, path)
print(f"\n=== FOMS RBAC 初始化完成 ===")
print(f" 机构: FOMS运维中心 ({ORG_ID})")
print(f" 角色: {', '.join(f'{k}({v})' for k, v in ROLES.items())}")
print(f" 用户:")
for u, i in USERS.items():
print(f" {u} / {i['password']} (角色: {i['role']})")
print(f"\n 角色权限分配:")
print(f" superadmin → 全部路径")
print(f" operator → 全部路径 (不含agent内部端点)")
print(f" viewer → 只读路径 (查看dashboard/列表)")
async def _add_roleperm(sor, role_id, perm_id, path):
"""Add rolepermission if not exists"""
rp_id = f'rp_{role_id}_{perm_id}'[:32]
existing = await sor.R('rolepermission', {'roleid': role_id, 'permid': perm_id})
if existing:
return
await sor.C('rolepermission', {
'id': rp_id,
'roleid': role_id,
'permid': perm_id,
})
if __name__ == '__main__':
asyncio.run(main())

View File

@ -1,97 +0,0 @@
#!/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",
# Metrics & logs
"/api/agent_report_metrics.dspy",
"/api/agent_push_logs.dspy",
"/api/node_detail.dspy",
"/api/node_metrics_history.dspy",
# Remote command
"/api/exec_remote_command.dspy",
"/api/agent_poll_commands.dspy",
"/api/agent_report_command_result.dspy",
]
# Agent heartbeat endpoint: no auth (agents use token)
PATHS_ANY.append("/api/agent_heartbeat.dspy")
PATHS_ANY.append("/api/agent_report_health.dspy")
PATHS_ANY.append("/api/agent_report_metrics.dspy")
PATHS_ANY.append("/api/agent_push_logs.dspy")
PATHS_ANY.append("/api/agent_poll_commands.dspy")
PATHS_ANY.append("/api/agent_report_command_result.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.")

View File

@ -1,6 +0,0 @@
#!/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}"

View File

@ -1,9 +0,0 @@
#!/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

View File

@ -1,7 +0,0 @@
# 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)

View File

@ -1,2 +0,0 @@
result = await agent_poll_commands(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,2 +0,0 @@
result = await agent_push_logs(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,2 +0,0 @@
result = await agent_report_command_result(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,3 +0,0 @@
# Agent health report endpoint (sync statuses)
result = await agent_report_health(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,2 +0,0 @@
result = await agent_report_metrics(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,3 +0,0 @@
# Dashboard stats API
stats = await get_dashboard_stats(request)
return json.dumps(stats, ensure_ascii=False, default=str)

View File

@ -1,2 +0,0 @@
result = await exec_remote_command(request, params_kw)
return json.dumps(result, ensure_ascii=False)

View File

@ -1,2 +0,0 @@
result = await get_node_detail(request, params_kw)
return json.dumps(result, ensure_ascii=False, default=str)

View File

@ -1,2 +0,0 @@
result = await get_node_metrics_history(request, params_kw)
return json.dumps(result, ensure_ascii=False, default=str)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,9 +0,0 @@
# Cluster create wrapper
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)

View File

@ -1,6 +0,0 @@
# 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)

View File

@ -1,6 +0,0 @@
# 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)

View File

@ -1,87 +0,0 @@
{
"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"}
}
]
}
]
}