diff --git a/app/foms.py b/app/foms.py index 4c93f95..d851cae 100644 --- a/app/foms.py +++ b/app/foms.py @@ -14,8 +14,11 @@ from ahserver.serverenv import ServerEnv from rbac.init import load_rbac from appbase.init import load_appbase -# Business modules +# Business modules(核心 + 三个拆出的独立模块) from foms.init import load_foms +from dbbackup.init import load_dbbackup +from filesync.init import load_filesync +from hostwatch.init import load_hostwatch from app.global_func import set_globalvariable @@ -25,7 +28,10 @@ def init(): load_pybricks() load_appbase() load_rbac() - load_foms() + load_foms() # 核心:集群/节点/切换/心跳/远程命令 + load_dbbackup() # 数据库备份 + 日志同步(binlog 复制) + load_filesync() # 运行文件同步(rsync) + load_hostwatch() # 日志监控 + 日志分析 + 主机指标 if __name__ == '__main__': diff --git a/build.sh b/build.sh index aff3149..4b3e58a 100755 --- a/build.sh +++ b/build.sh @@ -28,23 +28,50 @@ for m in appbase rbac; do cd $m && $cdir/py3/bin/pip install -e . done -# 5. Install FOMS +# 5. Install FOMS 核心 cd $cdir && $cdir/py3/bin/pip install -e . -# 6. Generate DDL +# 5.5 Install 业务模块(dbbackup / filesync / hostwatch) +for m in dbbackup filesync hostwatch; do + cd $cdir/pkgs + [ ! -d "$m" ] && git clone https://git.opencomputing.cn/yumoqing/$m + cd $m && $cdir/py3/bin/pip install -e . + # 生成模块 DDL + if [ -d models ]; then + json2ddl mysql models/ > models/mysql.ddl.sql 2>/dev/null || echo " skip ddl($m)" + fi +done + +# 6. Generate FOMS 核心 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/) +# 7. Generate FOMS 核心 CRUD UI (xls2ui for json/) cd $cdir if [ -d json ]; then cd json for f in *.json; do - echo " xls2ui $f" + echo " xls2ui foms $f" xls2ui -m ../models -o ../wwwroot foms "$f" 2>/dev/null || echo " skipped" done fi +# 7.5 生成模块 CRUD + 软链模块 wwwroot → 主 wwwroot/<模块名> +cd $cdir +for m in dbbackup filesync hostwatch; do + md=$cdir/pkgs/$m + # 生成模块 CRUD 页面 + if [ -d "$md/json" ]; then + cd "$md/json" + for f in *.json; do + echo " xls2ui $m $f" + xls2ui -m ../models -o ../wwwroot "$m" "$f" 2>/dev/null || echo " skipped" + done + fi + # 软链模块 wwwroot → 主 wwwroot/<模块名>(URL 前缀 /dbbackup/ 等) + ln -sfn "$md/wwwroot" "$cdir/wwwroot/$m" +done + # 8. Create runtime dirs mkdir -p $cdir/logs $cdir/files diff --git a/foms/init.py b/foms/init.py index 39d3059..74c8ee6 100644 --- a/foms/init.py +++ b/foms/init.py @@ -1,4 +1,20 @@ -"""FOMS business logic — CRUD operations + failover orchestration""" +"""FOMS 核心业务 — 主备切换编排(集群/节点/切换/心跳/远程命令)。 + +模块拆分后(2026-08-26),本模块只保留「主备切换」本身的核心概念与编排: + · 表:clusters、nodes、switchover_logs、node_commands + · 功能:集群/节点 CRUD、Agent 心跳、故障切换/切回、仪表盘、远程命令 + +拆出的独立模块(各自独立 git 仓库,表加模块名前缀): + · dbbackup — 数据库全量备份 + 日志同步(binlog 复制)管理 + 表 dbbackup_replications/jobs/records + · filesync — 运行文件同步(rsync) + 表 filesync_directories + · hostwatch — 日志监控 + 日志分析 + 主机指标 + 表 hostwatch_logs/metrics/rules/events + +跨模块引用采用松耦合:本模块只通过 cluster_id/node_id 字符串关联, +不 import 拆出模块的代码;仪表盘对拆出模块表做只读 SQL 聚合(同库)。 +""" import json from datetime import datetime from appPublic.uniqueID import getID @@ -9,6 +25,18 @@ from ahserver.serverenv import ServerEnv DBNAME = 'foms' +def _get(params_kw, key, default=''): + return params_kw.get(key, default) if hasattr(params_kw, 'get') else default + + +def _ns(params_kw): + ns = dict(params_kw) if hasattr(params_kw, 'items') else {} + for k in list(ns.keys()): + if k.endswith('_text'): + ns.pop(k) + return ns + + # ═══════════════════════════════════════════ # Cluster CRUD # ═══════════════════════════════════════════ @@ -16,15 +44,11 @@ DBNAME = 'foms' 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 = _ns(params_kw) 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'}} @@ -32,11 +56,8 @@ async def create_cluster(request, params_kw): 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 = _ns(params_kw) 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'}} @@ -44,8 +65,7 @@ async def update_cluster(request, params_kw): 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}) + await sor.D('clusters', {'id': _get(params_kw, 'id')}) return {'widgettype': 'Message', 'options': {'title': '成功', 'message': '集群已删除', 'type': 'success'}} @@ -56,14 +76,11 @@ async def delete_cluster(request, params_kw): 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 = _ns(params_kw) 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'}} @@ -71,11 +88,8 @@ async def create_node(request, params_kw): 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 = _ns(params_kw) 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'}} @@ -83,96 +97,18 @@ async def update_node(request, params_kw): 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}) + await sor.D('nodes', {'id': _get(params_kw, '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 +# Agent 心跳 # ═══════════════════════════════════════════ 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 '' + """Agent 心跳上报""" + node_id = _get(params_kw, 'node_id') + version = _get(params_kw, 'version') if not node_id: return {'status': 'error', 'message': 'node_id required'} env = ServerEnv() @@ -187,55 +123,22 @@ async def agent_heartbeat(request, params_kw): 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 '手动触发' + cluster_id = _get(params_kw, 'cluster_id') + reason = _get(params_kw, 'reason', '手动触发') 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, @@ -248,19 +151,14 @@ async def trigger_switchover(request, params_kw): '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, - } + 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 '' + cluster_id = _get(params_kw, 'cluster_id') log_id = getID() now = curDateString() async with get_sor_context(env, DBNAME) as sor: @@ -281,30 +179,29 @@ async def trigger_switchback(request, params_kw): 'created_at': now, }) await sor.U('clusters', {'id': cluster_id, 'status': 'switching', 'updated_at': now}) - return { - 'status': 'ok', - 'message': '切回流程已触发', - 'log_id': log_id, - } + return {'status': 'ok', 'message': '切回流程已触发', 'log_id': log_id} # ═══════════════════════════════════════════ -# Dashboard +# 仪表盘(跨模块只读聚合,同库 SQL) # ═══════════════════════════════════════════ 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'", {}) + # 拆出模块的表(同库松耦合只读) + dbs_synced = await sor.sqlExe( + "SELECT count(*) as cnt FROM dbbackup_replications WHERE sync_status='running'", {}) + dbs_error = await sor.sqlExe( + "SELECT count(*) as cnt FROM dbbackup_replications WHERE sync_status IN ('error','failed')", {}) + dirs_synced = await sor.sqlExe( + "SELECT count(*) as cnt FROM filesync_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', {} - ) + '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, @@ -322,192 +219,64 @@ async def get_dashboard_stats(request): # ═══════════════════════════════════════════ -# 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 +# 远程命令(专用 node_commands 表,不再借用日志表) # ═══════════════════════════════════════════ +# ⚠️ 安全待办:下发命令无鉴权白名单。管理中心被攻破 = 所有被管主机 root shell。 +# 后续独立任务应加:命令白名单 + RBAC 鉴权(仅 operator 以上可下发危险命令)。 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 + """下发远程命令,存入 node_commands 队列,Agent 轮询执行""" 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 '' + node_id = _get(params_kw, 'node_id') + command = _get(params_kw, 'command') 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', { + await sor.C('node_commands', { 'id': cmd_id, 'node_id': node_id, - 'log_source': 'foms_cmd', - 'log_level': 'INFO', - 'message': f'[PENDING] {command}', - 'logged_at': now, + 'command': command, + 'status': 'pending', + 'created_at': now, }) - return { - 'status': 'ok', - 'message': '命令已下发,等待Agent执行', - 'cmd_id': cmd_id, - } + return {'status': 'ok', 'message': '命令已下发,等待Agent执行', 'cmd_id': cmd_id} async def agent_poll_commands(request, params_kw): - """Agent轮询获取待执行命令""" + """Agent 轮询获取待执行命令""" env = ServerEnv() - node_id = params_kw.get('node_id', '') if hasattr(params_kw, 'get') else '' + node_id = _get(params_kw, 'node_id') 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} + "SELECT id, command FROM node_commands WHERE node_id=${nid}$ AND status='pending' " + "ORDER BY created_at ASC LIMIT 5", + {'nid': node_id}) + # 取到的命令标记为 running,防止重复下发 + for r in (rows or []): + await sor.U('node_commands', {'id': getattr(r, 'id', ''), 'status': 'running'}) + return {'commands': [{'id': getattr(r, 'id', ''), 'command': getattr(r, 'command', '')} + for r in (rows or [])]} async def agent_report_command_result(request, params_kw): - """Agent上报命令执行结果""" + """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 + cmd_id = _get(params_kw, 'cmd_id') async with get_sor_context(env, DBNAME) as sor: - await sor.U('node_logs', { + await sor.U('node_commands', { 'id': cmd_id, - 'log_source': 'foms_cmd', - 'log_level': 'INFO' if exit_code == 0 else 'ERROR', - 'message': f'[EXIT:{exit_code}] {output[:2000]}', + 'exit_code': _get(params_kw, 'exit_code', -1), + 'output': (_get(params_kw, 'output', '') or '')[:2000], + 'status': 'success' if int(_get(params_kw, 'exit_code', -1) or -1) == 0 else 'failed', + 'executed_at': curDateString(), }) return {'status': 'ok'} # ═══════════════════════════════════════════ -# ServerEnv Registration +# 注册 # ═══════════════════════════════════════════ def load_foms(): @@ -526,34 +295,14 @@ def load_foms(): 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 + # 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 diff --git a/json/sync_databases_list.json b/json/sync_databases_list.json deleted file mode 100644 index 7738b99..0000000 --- a/json/sync_databases_list.json +++ /dev/null @@ -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')}}" - } - } -} diff --git a/json/sync_directories_list.json b/json/sync_directories_list.json deleted file mode 100644 index 4401694..0000000 --- a/json/sync_directories_list.json +++ /dev/null @@ -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')}}" - } - } -} diff --git a/models/node_commands.json b/models/node_commands.json new file mode 100644 index 0000000..4f298c0 --- /dev/null +++ b/models/node_commands.json @@ -0,0 +1,27 @@ +{ + "summary": [ + { + "name": "node_commands", + "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": "command", "title": "命令内容", "type": "text", "nullable": "no"}, + {"name": "status", "title": "状态", "type": "str", "length": 16, "nullable": "no", "default": "pending"}, + {"name": "exit_code", "title": "退出码", "type": "int"}, + {"name": "output", "title": "执行输出", "type": "text"}, + {"name": "created_at", "title": "下发时间", "type": "datetime", "nullable": "no"}, + {"name": "executed_at", "title": "执行时间", "type": "datetime"} + ], + "indexes": [ + {"name": "idx_ncmd_node", "idxtype": "index", "idxfields": ["node_id"]}, + {"name": "idx_ncmd_status", "idxtype": "index", "idxfields": ["status"]} + ], + "codes": [ + {"field": "status", "table": "appcodes_kv", "valuefield": "k", "textfield": "v", "cond": "parentid='cmd_status'"} + ] +} diff --git a/models/node_logs.json b/models/node_logs.json deleted file mode 100644 index b6fe840..0000000 --- a/models/node_logs.json +++ /dev/null @@ -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'"} - ] -} diff --git a/models/node_metrics.json b/models/node_metrics.json deleted file mode 100644 index 37c99e8..0000000 --- a/models/node_metrics.json +++ /dev/null @@ -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"]} - ] -} diff --git a/models/sync_databases.json b/models/sync_databases.json deleted file mode 100644 index 360ce34..0000000 --- a/models/sync_databases.json +++ /dev/null @@ -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"} - ] -} diff --git a/models/sync_directories.json b/models/sync_directories.json deleted file mode 100644 index 37a86f9..0000000 --- a/models/sync_directories.json +++ /dev/null @@ -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"} - ] -} diff --git a/scripts/foms-agent.py b/scripts/foms-agent.py index 331ba61..b5a513d 100644 --- a/scripts/foms-agent.py +++ b/scripts/foms-agent.py @@ -56,7 +56,8 @@ class FomsAgent: # ═══ HTTP helpers ═══ def _api_get(self, path): - url = f"{MANAGEMENT_URL}/api/{path}" + # path 以 / 开头 = 完整路径(模块端点),否则走 /api/ 前缀(核心端点) + url = f"{MANAGEMENT_URL}{path}" if path.startswith('/') else f"{MANAGEMENT_URL}/api/{path}" req = urllib.request.Request(url) req.add_header('Authorization', f'Bearer {AGENT_TOKEN}') try: @@ -67,7 +68,7 @@ class FomsAgent: return None def _api_post(self, path, data): - url = f"{MANAGEMENT_URL}/api/{path}" + url = f"{MANAGEMENT_URL}{path}" if path.startswith('/') else 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}') @@ -202,28 +203,27 @@ class FomsAgent: return {'id': dir_config['id'], 'result': 'failed'} # ═══ Report Health ═══ - def report_health(self): - """Collect and report all health statuses""" + def report_db_health(self): + """收集 MySQL 复制状态 → 上报 dbbackup 模块""" 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) - + data = {'node_id': self.node_id, 'db_statuses': db_statuses} + self._api_post('/dbbackup/api/agent_report_db_health.dspy', data) + + def report_dir_health(self): + """执行 rsync 同步 → 上报 filesync 模块""" 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) + data = {'node_id': self.node_id, 'dir_statuses': dir_statuses} + self._api_post('/filesync/api/agent_report_dir_health.dspy', data) # ═══ Fetch Config from Management ═══ def fetch_config(self): @@ -254,9 +254,10 @@ class FomsAgent: self.heartbeat() now = time.time() - # Report health (MySQL status + rsync) + # Report health(MySQL 复制 → dbbackup;rsync → filesync) if now - last_health_report > HEALTH_REPORT_INTERVAL: - self.report_health() + self.report_db_health() + self.report_dir_health() last_health_report = now # Collect & report system metrics @@ -314,7 +315,7 @@ class FomsAgent: if not metrics: return data = {'node_id': self.node_id, 'metrics': metrics} - result = self._api_post('agent_report_metrics.dspy', data) + result = self._api_post('/hostwatch/api/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: diff --git a/wwwroot/api/agent_push_logs.dspy b/wwwroot/api/agent_push_logs.dspy deleted file mode 100644 index 38a480a..0000000 --- a/wwwroot/api/agent_push_logs.dspy +++ /dev/null @@ -1,2 +0,0 @@ -result = await agent_push_logs(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 deleted file mode 100644 index b8e1561..0000000 --- a/wwwroot/api/agent_report_health.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/agent_report_metrics.dspy b/wwwroot/api/agent_report_metrics.dspy deleted file mode 100644 index 2a76301..0000000 --- a/wwwroot/api/agent_report_metrics.dspy +++ /dev/null @@ -1,2 +0,0 @@ -result = await agent_report_metrics(request, params_kw) -return json.dumps(result, ensure_ascii=False) diff --git a/wwwroot/api/node_detail.dspy b/wwwroot/api/node_detail.dspy deleted file mode 100644 index 4970bc3..0000000 --- a/wwwroot/api/node_detail.dspy +++ /dev/null @@ -1,2 +0,0 @@ -result = await get_node_detail(request, params_kw) -return json.dumps(result, ensure_ascii=False, default=str) diff --git a/wwwroot/api/node_metrics_history.dspy b/wwwroot/api/node_metrics_history.dspy deleted file mode 100644 index ab16f77..0000000 --- a/wwwroot/api/node_metrics_history.dspy +++ /dev/null @@ -1,2 +0,0 @@ -result = await get_node_metrics_history(request, params_kw) -return json.dumps(result, ensure_ascii=False, default=str) diff --git a/wwwroot/api/sync_databases_create.dspy b/wwwroot/api/sync_databases_create.dspy deleted file mode 100644 index 7007c18..0000000 --- a/wwwroot/api/sync_databases_create.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/sync_databases_delete.dspy b/wwwroot/api/sync_databases_delete.dspy deleted file mode 100644 index 90c8b64..0000000 --- a/wwwroot/api/sync_databases_delete.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/sync_databases_update.dspy b/wwwroot/api/sync_databases_update.dspy deleted file mode 100644 index d1da541..0000000 --- a/wwwroot/api/sync_databases_update.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/sync_directories_create.dspy b/wwwroot/api/sync_directories_create.dspy deleted file mode 100644 index 358d08c..0000000 --- a/wwwroot/api/sync_directories_create.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/sync_directories_delete.dspy b/wwwroot/api/sync_directories_delete.dspy deleted file mode 100644 index 95c83ba..0000000 --- a/wwwroot/api/sync_directories_delete.dspy +++ /dev/null @@ -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) diff --git a/wwwroot/api/sync_directories_update.dspy b/wwwroot/api/sync_directories_update.dspy deleted file mode 100644 index 8d44f9a..0000000 --- a/wwwroot/api/sync_directories_update.dspy +++ /dev/null @@ -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)