提交 27c461c2 authored 作者: 陈泽健's avatar 陈泽健

feat(service_monitor): 61_platform_log_check 平台日志分析模块上线(含遗留变更一并入库)

61 平台日志分析模块(仅新统一平台 /data/services 生效):
- 新增 assets/system/61_platform_log_check.sh:18 服务三组检测
  (core 10 完整含活跃性 / dubbo 2 只检文件+ERROR / missing 6 固定严重),
  ERROR 密度 + 最近 ERROR 时间节点 + 汇总行,KEY 规范 LOGCHK_<大写>_<项> 且每项独立 _LEVEL
- config.sh.template:服务清单 3 组 + 阈值 5 键(静默 15/60min、密度 50/100、窗口 2000 行)
- check_modules 注册(仅 full 套件)+ 阈值白名单;display_names 130 键;测试 +4 用例

一并入库此前遗留(HANDOFF §7 应提交清单):
- logger 迁移至 utils.logger(compare/notification/report/runner/statistics/target_service + check_modules)
- executor 上传强制 LF + assets/*.sh CRLF 修复
- report_service check_missing 仅查定时目标;runner_service 通知留痕 + 僵尸运行清理
- tests/conftest.py 同步

其他:
- HANDOFF §16 记录本模块设计/踩坑/验证
- .gitignore 追加 data/report_index.json(运行时数据勿入库)

验证:5.44 真实机 1.5s/106 键全 PASS;5.60 部署后旧平台短路正确;全量回归 94 用例通过。
Co-Authored-By: 's avatarClaude Code <noreply@anthropic.com>
上级 80b76ee3
...@@ -62,6 +62,7 @@ temp/ ...@@ -62,6 +62,7 @@ temp/
# 服务监测模块运行时数据(含加密凭据与报告,勿入库) # 服务监测模块运行时数据(含加密凭据与报告,勿入库)
skill/code/web/service_monitor/data/targets.json skill/code/web/service_monitor/data/targets.json
skill/code/web/service_monitor/data/reports/*.json skill/code/web/service_monitor/data/reports/*.json
skill/code/web/service_monitor/data/report_index.json
!skill/code/web/service_monitor/data/.gitkeep !skill/code/web/service_monitor/data/.gitkeep
!skill/code/web/service_monitor/data/reports/.gitkeep !skill/code/web/service_monitor/data/reports/.gitkeep
......
# HANDOFF — 服务监测模块实施进度 # HANDOFF — 服务监测模块实施进度
> 最后更新:2026-09-01 | 分支:troubleshoot-ai-assistant | 模块:service-monitor > 最后更新:2026-09-09 | 分支:troubleshoot-ai-assistant | 模块:service-monitor
> 状态:**5.202 定时任务 end_date 过期问题已修复并部署 5.60(7249d355);5.202 end_date 已延长至 2027-08-31 恢复巡检;HANDOFF 已精简为运维速查(06369d48)** > 状态:**61_platform_log_check 平台日志分析模块已上线并部署 5.60(9fb743cc);仅新平台生效,18 服务三组检测**
> 📦 完整历史会话进度见归档:`Docs/需求文档/服务监测/HANDOFF_归档_2026-09-01.md` > 📦 完整历史会话进度见归档:`Docs/需求文档/服务监测/HANDOFF_归档_2026-09-01.md`
> (本文仅保留运维必需信息 + 当前状态;历史根因/部署细节已在归档中按会话编号可查) > (本文仅保留运维必需信息 + 当前状态;历史根因/部署细节已在归档中按会话编号可查)
...@@ -109,6 +109,7 @@ skill/code/web/service_monitor/ ...@@ -109,6 +109,7 @@ skill/code/web/service_monitor/
| 13 | 09-01 | 提交 HANDOFF 更新 | **eebe75d4** | | 13 | 09-01 | 提交 HANDOFF 更新 | **eebe75d4** |
| 14 | 09-01 | **5.202 end_date 过期修复 + 部署 5.60 + 运维延长恢复**(本次,详见下) | **7249d355** | | 14 | 09-01 | **5.202 end_date 过期修复 + 部署 5.60 + 运维延长恢复**(本次,详见下) | **7249d355** |
| 15 | 09-01 | **HANDOFF 精简 677→146 行 + 归档完整历史**(本文件) | **06369d48** | | 15 | 09-01 | **HANDOFF 精简 677→146 行 + 归档完整历史**(本文件) | **06369d48** |
| 16 | 09-08~09 | **61_platform_log_check 平台日志分析模块:18 服务三组检测(ERROR 密度+时间节点+活跃性)、仅新平台生效;5.44 实测 + 5.60 部署**(详见下) | **9fb743cc** |
### §14 本次:5.202 定时任务 end_date 过期(修复 + 部署 + 运维恢复) ### §14 本次:5.202 定时任务 end_date 过期(修复 + 部署 + 运维恢复)
...@@ -127,16 +128,44 @@ skill/code/web/service_monitor/ ...@@ -127,16 +128,44 @@ skill/code/web/service_monitor/
--- ---
### §16 本次:61_platform_log_check 平台日志分析模块(新增 + 5.44 实测 + 5.60 部署)
**需求**:现有日志类检测(33 数 ERROR 文件数 / 44 查 /var/log / 59 查存在性大小)未覆盖 `/data/services/api/` 下业务服务日志质量。用户需求“主要看有没有 ERROR,有则记录时间节点”。
**模块设计(9fb743cc,已推送)**
- 新增 `assets/system/61_platform_log_check.sh`**仅新平台**`/data/services` 存在)生效;旧平台输出 INFO 跳过
- 18 服务三组:**core 10**(auth×3+meeting×5+python)完整检测(文件/ERROR/活跃性);**dubbo 2**(teams/meeting-control)只检文件/ERROR;**missing 6**(ews-cloud/serviceCall/zoom/cloudLink/smc-three/xylink)预期缺失每轮固定严重
- 每服务输出 `LOGCHK_<大写下划线>_FILE/ERROR_CNT/ERROR_LAST` 等,**每个检测项带独立 `<KEY>_LEVEL`**,供 parser 精确关联状态
- 汇总项 `PLATFORM_LOG_SUMMARY`(共18个 正常N 警告N 严重N)+ `_SUMMARY_LEVEL`
**KEY 规范(parser 兼容,本轮最大踩坑)**
- parser `_KV_RE` 只认大写/数字/下划线 KEY,小写/连字符/点号会被静默丢弃 → 前缀必须 `LOGCHK_` + 服务名 `tr a-z→A-Z | tr .-→__`(如 `java-meeting2.0``LOGCHK_JAVA_MEETING2_0`
- parser 跳过 `_LEVEL` 结尾键且状态只做**精确** `<KEY>_LEVEL` 关联 → 不能只用一个全局 `<P>_LEVEL`
**bash 陷阱(已修)**
- `grep -c ... || echo "0"` 无匹配时输出 "0\n0" 双行 + `[: 需要整数表达式` → 改 `|| true`
- `grep 'ERROR' || true | tail -1` 优先级陷阱(grep 成功时 tail 不执行)→ 改 `grep 'ERROR' | tail -1`
**验证**
- 5.44 真实机探针(`.tmp/`,明文密码不入库):1.5s、106 键、12 存在含 ERROR_LAST、6 缺失严重;`SUMMARY: 共18 正常8 警告1 严重9 | 严重`;抽样 python-voice 357 严重 / dubbo-teams 170 严重(09-07 04:07)/ java-quartz 21 正常
- 5.60 部署:文件就位、0 CRLF、bash -n OK、Python 加载 module=平台日志分析 130 显示键、旧平台短路正确
- 全量回归 94 用例通过
**文档**`PRD_计划执行_新统一平台日志分析监测模块.md`
---
## 7. 工作区遗留未提交(下次提交范围参考) ## 7. 工作区遗留未提交(下次提交范围参考)
**✅ 应提交(服务监测正式代码,已测试):** **✅ 已提交(9fb743cc,09-09):**
- logger 迁移:`services/{compare,notification,report,runner,statistics,target}_service.py``utils/check_modules.py`(→ `utils.logger.get_logger` - 61_platform_log_check 模块全套(脚本/config 模板/check_modules/display_names/测试)
- CRLF:`utils/executor.py`(强制 LF)+ `assets/*.sh`/`config.sh.template` - logger 迁移(6 个 service + check_modules)、executor CRLF 强制 LF、assets/*.sh CRLF
- 功能:`report_service.py`(check_missing 仅查定时目标)、`runner_service.py`(通知结果日志 + 僵尸清理)、`schedule_service.py`(已随 7249d355 提交) - `report_service.py`(check_missing 仅查定时目标)、`runner_service.py`(通知结果日志 + 僵尸清理)、`tests/conftest.py`
- 测试:`tests/conftest.py`
**✅ 应提交(待办):**
- 服务管理:`services/five44_client.py``services/service_manage.py`(另立 commit) - 服务管理:`services/five44_client.py``services/service_manage.py`(另立 commit)
**⚠️ 其他模块/勿提交:** `Dockerfile``.gitattributes``Docs/服务管理/``HANDOFF_容器部署_5.69.md``frontend/tsconfig.app.tsbuildinfo`(build 产物)、`users.json`(运行时变更)、`deploy/*.py` 运维脚本(已确认保留,是否入库视需要)、`Docs/维护手册/`(111MB 勿入库)、`skill/code/搜索索引.json`/`搜索向量.json`(运行时生成) **⚠️ 其他模块/勿提交:** `Dockerfile``.gitattributes``Docs/服务管理/``HANDOFF_容器部署_5.69.md``frontend/tsconfig.app.tsbuildinfo`(build 产物)、`users.json`(运行时变更)、`deploy/*.py` 运维脚本(已确认保留,是否入库视需要)、`Docs/维护手册/`(111MB 勿入库)、`skill/code/搜索索引.json`/`搜索向量.json`(运行时生成)`service_monitor/data/report_index.json`(已 gitignore)
--- ---
......
...@@ -144,3 +144,24 @@ UPYTHON_PORTS_OLD="${UPYTHON_PORTS_OLD:-8000:uwsgi,11211:memcached,8081:nginx,84 ...@@ -144,3 +144,24 @@ UPYTHON_PORTS_OLD="${UPYTHON_PORTS_OLD:-8000:uwsgi,11211:memcached,8081:nginx,84
# ==================== 配置IP检测配置 ==================== # ==================== 配置IP检测配置 ====================
# 服务器IP(运行时注入,用于IP允许列表) # 服务器IP(运行时注入,用于IP允许列表)
SERVER_IP="${SERVER_IP:-}" SERVER_IP="${SERVER_IP:-}"
# ==================== 平台日志分析配置(61_platform_log_check.sh) ====================
# 服务清单格式: prefix:相对/data/services/api的日志路径,prefix:路径,...(3 组:core 完整检测 / dubbo 只检文件+ERROR / missing 预期缺失)
PLATFORM_LOG_SERVICES_CORE="${PLATFORM_LOG_SERVICES_CORE:-auth-sso-auth:auth/auth-sso-auth/logs/ubains-INFO-AND-ERROR.log,auth-sso-gatway:auth/auth-sso-gatway/logs/ubains-INFO-AND-ERROR.log,auth-sso-system:auth/auth-sso-system/logs/ubains-INFO-AND-ERROR.log,java-meeting2.0:java-meeting/java-meeting2.0/logs/ubains-INFO-AND-ERROR.log,java-meeting-extapi:java-meeting/java-meeting-extapi/logs/ubains-INFO-AND-ERROR.log,java-message-scheduling:java-meeting/java-message-scheduling/logs/ubains-INFO-AND-ERROR.log,java-mqtt:java-meeting/java-mqtt/logs/ubains-INFO-AND-ERROR.log,java-quartz:java-meeting/java-quartz/logs/ubains-INFO-AND-ERROR.log,python-cmdb:python-cmdb/log/uinfo.log,python-voice:python-voice/log/uinfo.log}"
# dubbo 服务(2 个有主日志,只检文件/体积/ERROR,不做活跃性)
PLATFORM_LOG_SERVICES_DUBBO="${PLATFORM_LOG_SERVICES_DUBBO:-dubbo-teams:dubbo/dubbo-teams/logs/ubains-INFO-AND-ERROR.log,dubbo-meeting-control:dubbo/dubbo-meeting-control/logs/ubains-INFO-AND-ERROR.log}"
# 日志缺失服务(6 个,已部署但无主日志,每轮固定报严重)
PLATFORM_LOG_SERVICES_MISSING="${PLATFORM_LOG_SERVICES_MISSING:-dubbo-ews-cloud:dubbo/dubbo-ews-cloud/logs/ubains-INFO-AND-ERROR.log,dubbo-serviceCall:dubbo/dubbo-serviceCall/logs/ubains-INFO-AND-ERROR.log,dubbo-zoom:dubbo/dubbo-zoom/logs/ubains-INFO-AND-ERROR.log,dubbo-cloudLink:dubbo/dubbo-cloudLink/logs/ubains-INFO-AND-ERROR.log,dubbo-smc-three:dubbo/dubbo-smc-three/logs/ubains-INFO-AND-ERROR.log,dubbo-xylink:dubbo/dubbo-xylink/logs/ubains-INFO-AND-ERROR.log}"
# 日志静默失联阈值(分钟,仅 core 组生效)
PLATFORM_LOG_STALE_WARN_MIN="${PLATFORM_LOG_STALE_WARN_MIN:-15}"
PLATFORM_LOG_STALE_CRIT_MIN="${PLATFORM_LOG_STALE_CRIT_MIN:-60}"
# ERROR 密度阈值(最近 N 行内 ERROR 条数)
PLATFORM_LOG_ERR_CNT_WARN="${PLATFORM_LOG_ERR_CNT_WARN:-50}"
PLATFORM_LOG_ERR_CNT_CRIT="${PLATFORM_LOG_ERR_CNT_CRIT:-100}"
# 扫描窗口行数(tail 行数)
PLATFORM_LOG_SCAN_LINES="${PLATFORM_LOG_SCAN_LINES:-2000}"
#!/bin/bash
################################################################################
# 61_platform_log_check.sh — 平台日志分析检测(仅新统一平台)
# 功能: 检测 /data/services/api/ 下各业务服务日志的存在性/活跃性/ERROR密度/最近ERROR时间节点
# 设计: 三组服务
# - core 核心服务(10): auth×3 + meeting×5 + python-cmdb + python-voice,做完整检测(含活跃性)
# - dubbo dubbo服务(2): dubbo-teams / dubbo-meeting-control,只检文件/ERROR(日志低频更新属正常,不做活跃性)
# - missing 日志缺失服务(6): 已部署但无主日志,每轮固定报严重
# 平台门: 仅 /data/services 存在(新平台)时生效;旧平台不输出任何检测项
# KEY 规范: 输出键统一 LOGCHK_<服务名大写下划线>_<项>,且每个检测项带 <KEY>_LEVEL
# (parser 按精确 <KEY>_LEVEL 关联状态,无则状态回落到默认正常)
# 日期: 2026-09-08
################################################################################
# 获取脚本所在目录并加载依赖
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
LIB_DIR="${LIB_DIR:-/tmp/check_modules}"
# 加载配置文件和通用函数库
if [ -f "$LIB_DIR/lib/config.sh" ]; then
source "$LIB_DIR/lib/config.sh"
else
echo "ERROR: 配置文件不存在: $LIB_DIR/lib/config.sh"
exit 1
fi
if [ -f "$LIB_DIR/lib/common.sh" ]; then
source "$LIB_DIR/lib/common.sh"
else
echo "ERROR: 通用函数库不存在: $LIB_DIR/lib/common.sh"
exit 1
fi
# ==================== 平台类型识别 ====================
detect_platform_type() {
[ -d "/data/services" ] && echo "new" || echo "old"
}
# ==================== 服务名 → LOGCHK 键前缀 ====================
# 例: auth-sso-auth → LOGCHK_AUTH_SSO_AUTH;java-meeting2.0 → LOGCHK_JAVA_MEETING2_0
to_key_prefix() {
echo "LOGCHK_$(echo "$1" | tr 'a-z' 'A-Z' | tr '.-' '__')"
}
# ==================== 时间戳提取 ====================
# 兼容两种格式:
# Java: 2026-09-08 11:11:01.105 [http-nio-...] ERROR ...
# Python:[2026-09-07 18:12:05,483][...][ERROR]...
# 取行首 YYYY-MM-DD HH:MM
extract_timestamp() {
local line="$1"
local ts
ts=$(echo "$line" | sed -nE 's/^\[?([0-9]{4}-[0-9]{2}-[0-9]{2} [0-9]{2}:[0-9]{2}).*/\1/p')
[ -n "$ts" ] && echo "$ts" || echo "无"
}
# ==================== 状态计数(汇总用) ====================
_PLOG_TOTAL=0
_PLOG_OK=0
_PLOG_WARN=0
_PLOG_CRIT=0
plog_count_level() {
local level="$1"
_PLOG_TOTAL=$((_PLOG_TOTAL + 1))
case "$level" in
"严重") _PLOG_CRIT=$((_PLOG_CRIT + 1)) ;;
"警告") _PLOG_WARN=$((_PLOG_WARN + 1)) ;;
*) _PLOG_OK=$((_PLOG_OK + 1)) ;;
esac
}
# 取两个等级中更严重者(用于每项状态)
worst_level() {
local a="$1" b="$2"
if [ "$a" = "严重" ] || [ "$b" = "严重" ]; then
echo "严重"
elif [ "$a" = "警告" ] || [ "$b" = "警告" ]; then
echo "警告"
else
echo "正常"
fi
}
# ==================== 单服务检测 ====================
# 参数: prefix relpath check_mtime(0/1)
# 输出(含文件缺失分叉):
# 缺失: <P>_FILE=缺失 + <P>_FILE_LEVEL=严重
# 存在: <P>_FILE=存在 + <P>_FILE_LEVEL=正常
# [<P>_MTIME_MIN=N + <P>_MTIME_MIN_LEVEL=...] (仅 core)
# <P>_ERROR_CNT=N + <P>_ERROR_CNT_LEVEL=...
# <P>_ERROR_LAST=时间/无 + <P>_ERROR_LAST_LEVEL=...
# 每项均带精确 <KEY>_LEVEL,供 parser 直接关联状态。
check_one_service() {
local prefix="$1"
local relpath="$2"
local check_mtime="$3"
local kp
kp=$(to_key_prefix "$prefix")
local file="/data/services/api/$relpath"
if [ ! -f "$file" ]; then
output_result "${kp}_FILE" "缺失"
output_result "${kp}_FILE_LEVEL" "严重"
plog_count_level "严重"
return
fi
output_result "${kp}_FILE" "存在"
output_result "${kp}_FILE_LEVEL" "正常"
# 活跃性(仅 core 组)
local mtime_min=0 mtime_level="正常"
if [ "$check_mtime" = "1" ]; then
local mtime now
mtime=$(stat -c %Y "$file" 2>/dev/null || echo "0")
now=$(date +%s)
mtime_min=$(( (now - mtime) / 60 ))
[ "$mtime_min" -lt 0 ] && mtime_min=0
output_result "${kp}_MTIME_MIN" "$mtime_min"
local stale_crit stale_warn
stale_crit="${PLATFORM_LOG_STALE_CRIT_MIN:-60}"
stale_warn="${PLATFORM_LOG_STALE_WARN_MIN:-15}"
if [ "$mtime_min" -gt "$stale_crit" ]; then
mtime_level="严重"
elif [ "$mtime_min" -gt "$stale_warn" ]; then
mtime_level="警告"
fi
output_result "${kp}_MTIME_MIN_LEVEL" "$mtime_level"
fi
# ERROR 密度 + 最近 ERROR 时间(最近 N 行窗口)
# 注意: grep 无匹配退出码为 1,用 || true 兜底避免命令替换输出脏值
local scan_lines scan_tail err_cnt err_last_line err_last
scan_lines="${PLATFORM_LOG_SCAN_LINES:-2000}"
scan_tail=$(tail -n "$scan_lines" "$file" 2>/dev/null || true)
err_cnt=$(printf '%s\n' "$scan_tail" | grep -c 'ERROR' || true)
err_last_line=$(printf '%s\n' "$scan_tail" | grep 'ERROR' | tail -1)
if [ -n "$err_last_line" ]; then
err_last=$(extract_timestamp "$err_last_line")
else
err_last="无"
fi
output_result "${kp}_ERROR_CNT" "$err_cnt"
output_result "${kp}_ERROR_LAST" "$err_last"
local err_level="正常"
local err_crit err_warn
err_crit="${PLATFORM_LOG_ERR_CNT_CRIT:-100}"
err_warn="${PLATFORM_LOG_ERR_CNT_WARN:-50}"
if [ "$err_cnt" -gt "$err_crit" ]; then
err_level="严重"
elif [ "$err_cnt" -gt "$err_warn" ]; then
err_level="警告"
fi
output_result "${kp}_ERROR_CNT_LEVEL" "$err_level"
output_result "${kp}_ERROR_LAST_LEVEL" "$err_level"
# 服务整体等级(活跃性 + ERROR 密度取较严重者),仅用于汇总计数
local overall
overall=$(worst_level "$mtime_level" "$err_level")
plog_count_level "$overall"
}
# ==================== 服务列表遍历 ====================
# 参数: services_str("prefix:relpath,prefix:relpath") check_mtime(0/1)
check_service_list() {
local services_str="$1"
local check_mtime="$2"
local IFS_save="$IFS"
IFS=','
local services=($services_str)
IFS="$IFS_save"
for svc in "${services[@]}"; do
[ -z "$svc" ] && continue
local prefix relpath
prefix="${svc%%:*}"
relpath="${svc#*:}"
[ -z "$prefix" ] || [ -z "$relpath" ] && continue
check_one_service "$prefix" "$relpath" "$check_mtime"
done
}
# ==================== 主检测流程 ====================
main() {
log_info "开始平台日志分析..."
set +e
local platform_type
platform_type=$(detect_platform_type)
if [ "$platform_type" != "new" ]; then
log_info "非新统一平台,跳过平台日志分析"
return
fi
# 三组服务清单(配置驱动,模板默认值)
local core_str dubbo_str missing_str
core_str="${PLATFORM_LOG_SERVICES_CORE:-auth-sso-auth:auth/auth-sso-auth/logs/ubains-INFO-AND-ERROR.log,auth-sso-gatway:auth/auth-sso-gatway/logs/ubains-INFO-AND-ERROR.log,auth-sso-system:auth/auth-sso-system/logs/ubains-INFO-AND-ERROR.log,java-meeting2.0:java-meeting/java-meeting2.0/logs/ubains-INFO-AND-ERROR.log,java-meeting-extapi:java-meeting/java-meeting-extapi/logs/ubains-INFO-AND-ERROR.log,java-message-scheduling:java-meeting/java-message-scheduling/logs/ubains-INFO-AND-ERROR.log,java-mqtt:java-meeting/java-mqtt/logs/ubains-INFO-AND-ERROR.log,java-quartz:java-meeting/java-quartz/logs/ubains-INFO-AND-ERROR.log,python-cmdb:python-cmdb/log/uinfo.log,python-voice:python-voice/log/uinfo.log}"
dubbo_str="${PLATFORM_LOG_SERVICES_DUBBO:-dubbo-teams:dubbo/dubbo-teams/logs/ubains-INFO-AND-ERROR.log,dubbo-meeting-control:dubbo/dubbo-meeting-control/logs/ubains-INFO-AND-ERROR.log}"
missing_str="${PLATFORM_LOG_SERVICES_MISSING:-dubbo-ews-cloud:dubbo/dubbo-ews-cloud/logs/ubains-INFO-AND-ERROR.log,dubbo-serviceCall:dubbo/dubbo-serviceCall/logs/ubains-INFO-AND-ERROR.log,dubbo-zoom:dubbo/dubbo-zoom/logs/ubains-INFO-AND-ERROR.log,dubbo-cloudLink:dubbo/dubbo-cloudLink/logs/ubains-INFO-AND-ERROR.log,dubbo-smc-three:dubbo/dubbo-smc-three/logs/ubains-INFO-AND-ERROR.log,dubbo-xylink:dubbo/dubbo-xylink/logs/ubains-INFO-AND-ERROR.log}"
check_service_list "$core_str" "1"
check_service_list "$dubbo_str" "0"
check_service_list "$missing_str" "0"
# 汇总项(单行,parser 经 PLATFORM_LOG_SUMMARY_LEVEL 关联状态)
local plevel="正常"
[ "$_PLOG_CRIT" -gt 0 ] && plevel="严重"
[ "$plevel" = "正常" ] && [ "$_PLOG_WARN" -gt 0 ] && plevel="警告"
output_result "PLATFORM_LOG_SUMMARY" "共${_PLOG_TOTAL}个服务 正常${_PLOG_OK} 警告${_PLOG_WARN} 严重${_PLOG_CRIT}"
output_result "PLATFORM_LOG_SUMMARY_LEVEL" "$plevel"
log_info "平台日志分析完成"
}
main
\ No newline at end of file
...@@ -9,13 +9,13 @@ compare_service.py — 报告对比服务 ...@@ -9,13 +9,13 @@ compare_service.py — 报告对比服务
from __future__ import annotations from __future__ import annotations
import logging
import re import re
from typing import Dict, List, Optional, Tuple from typing import Dict, List, Optional, Tuple
from .report_service import get_report, NORMAL, WARNING, CRITICAL from .report_service import get_report, NORMAL, WARNING, CRITICAL
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.compare_service") logger = get_logger("service_monitor.compare_service")
# 值变化阈值(变化率 > 20% 视为"值变化") # 值变化阈值(变化率 > 20% 视为"值变化")
VALUE_CHANGE_THRESHOLD = 0.2 VALUE_CHANGE_THRESHOLD = 0.2
......
...@@ -15,7 +15,6 @@ from __future__ import annotations ...@@ -15,7 +15,6 @@ from __future__ import annotations
import hashlib import hashlib
import hmac import hmac
import json import json
import logging
import smtplib import smtplib
import time import time
import urllib.parse import urllib.parse
...@@ -26,8 +25,9 @@ from typing import Optional ...@@ -26,8 +25,9 @@ from typing import Optional
from ..utils.paths import NOTIFICATIONS_FILE, ensure_dirs from ..utils.paths import NOTIFICATIONS_FILE, ensure_dirs
from ..utils.crypto import encrypt_password, decrypt_password from ..utils.crypto import encrypt_password, decrypt_password
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.notification") logger = get_logger("service_monitor.notification")
# 默认配置模板 # 默认配置模板
DEFAULT_CONFIG = { DEFAULT_CONFIG = {
......
...@@ -9,7 +9,6 @@ report_service.py — 巡检报告管理 ...@@ -9,7 +9,6 @@ report_service.py — 巡检报告管理
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import secrets import secrets
import threading import threading
import uuid import uuid
...@@ -23,8 +22,9 @@ except ImportError: ...@@ -23,8 +22,9 @@ except ImportError:
SHANGHAI_TZ = None SHANGHAI_TZ = None
from ..utils.paths import REPORTS_DIR, COOLDOWN_FILE, REPORT_INDEX_FILE, ensure_dirs from ..utils.paths import REPORTS_DIR, COOLDOWN_FILE, REPORT_INDEX_FILE, ensure_dirs
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.report_service") logger = get_logger("service_monitor.report_service")
# 报告默认保留天数(可被配置覆盖) # 报告默认保留天数(可被配置覆盖)
DEFAULT_RETENTION_DAYS = 14 DEFAULT_RETENTION_DAYS = 14
...@@ -841,6 +841,11 @@ def check_missing_reports(missing_days: int = 2) -> list: ...@@ -841,6 +841,11 @@ def check_missing_reports(missing_days: int = 2) -> list:
对每个目标查询最新报告时间,如果超过 missing_days 天未生成报告, 对每个目标查询最新报告时间,如果超过 missing_days 天未生成报告,
则返回告警列表。 则返回告警列表。
仅检查配置了启用定时任务的目标:内置目标(如"本机(当前服务器)")
若从未配置定时巡检,不参与报告缺失判定,避免每天误报
"XX 天无新报告(上次:无报告)"。schedules.json 读取失败时回退为
检查全部目标(保持原行为兜底)。
Args: Args:
missing_days: 缺失天数阈值(默认 2 天) missing_days: 缺失天数阈值(默认 2 天)
...@@ -855,7 +860,18 @@ def check_missing_reports(missing_days: int = 2) -> list: ...@@ -855,7 +860,18 @@ def check_missing_reports(missing_days: int = 2) -> list:
... ...
] ]
""" """
from . import target_service from . import target_service, schedule_service
# 收集配置了启用定时任务的目标 ID;读不到调度配置时置 None(检查全部)
scheduled_ids = None
try:
scheduled_ids = {
s.get("target_id")
for s in schedule_service.list_schedules()
if s.get("enabled") and s.get("target_id")
}
except Exception:
logger.warning("读取定时任务列表失败,回退为检查全部目标", exc_info=True)
missing = [] missing = []
now = datetime.now() now = datetime.now()
...@@ -865,6 +881,10 @@ def check_missing_reports(missing_days: int = 2) -> list: ...@@ -865,6 +881,10 @@ def check_missing_reports(missing_days: int = 2) -> list:
target_id = target.get("id", "") target_id = target.get("id", "")
target_name = target.get("name", "未知") target_name = target.get("name", "未知")
# 无任何启用定时任务的目标不参与报告缺失判定(如内置"本机")
if scheduled_ids is not None and target_id not in scheduled_ids:
continue
# 获取最新一份报告 # 获取最新一份报告
reports = list_reports(target_id=target_id, limit=1) reports = list_reports(target_id=target_id, limit=1)
if not reports: if not reports:
......
...@@ -15,7 +15,6 @@ from __future__ import annotations ...@@ -15,7 +15,6 @@ from __future__ import annotations
import hashlib import hashlib
import hmac import hmac
import logging
import threading import threading
import time import time
import traceback import traceback
...@@ -37,6 +36,10 @@ logger = get_logger("service_monitor.runner") ...@@ -37,6 +36,10 @@ logger = get_logger("service_monitor.runner")
_runs: Dict[str, dict] = {} _runs: Dict[str, dict] = {}
_runs_lock = threading.Lock() _runs_lock = threading.Lock()
# 僵尸巡检超时阈值(秒):超过此时长未完成的 run 视为僵尸自动清除。
# full 套件 42 模块正常 2-3 分钟完成,10 分钟仍未结束必为卡死。
_RUN_STALE_TIMEOUT = 600
def _now() -> str: def _now() -> str:
return datetime.now().strftime("%Y-%m-%dT%H:%M:%S") return datetime.now().strftime("%Y-%m-%dT%H:%M:%S")
...@@ -53,6 +56,7 @@ def _register_run(run_id: str, target_id: str, suite: str, total: int) -> None: ...@@ -53,6 +56,7 @@ def _register_run(run_id: str, target_id: str, suite: str, total: int) -> None:
"cancel": threading.Event(), "cancel": threading.Event(),
"finished": False, "finished": False,
"error": None, "error": None,
"started_ts": time.time(),
} }
...@@ -67,6 +71,37 @@ def _pop_run(run_id: str) -> Optional[dict]: ...@@ -67,6 +71,37 @@ def _pop_run(run_id: str) -> Optional[dict]:
return _runs.get(run_id) return _runs.get(run_id)
def _cleanup_stale_runs() -> None:
"""清除僵尸巡检:超时且未完成的 run 从注册表中移除。
场景:paramiko recv_exit_status() 无限阻塞导致巡检线程卡死,
_runs 中遗留未完成记录,使该目标后续所有巡检命中防重复逻辑返回 409,
前端表现为永久“连接中断”。此处在查询运行状态前兜底清理,
配合 executor 的 channel 超时关闭双保险。
"""
now = time.time()
stale_ids = []
with _runs_lock:
for run_id, run in _runs.items():
if run.get("finished"):
continue
started_ts = run.get("started_ts")
if started_ts and (now - started_ts) > _RUN_STALE_TIMEOUT:
stale_ids.append(run_id)
for run_id in stale_ids:
run = _runs.pop(run_id, None)
logger.warning(
"清除僵尸巡检 run_id=%s target=%s suite=%s done=%s/%s "
"(超过 %ds 未完成)",
run_id,
run.get("target_id") if run else "?",
run.get("suite") if run else "?",
run.get("done", 0) if run else 0,
run.get("total", 0) if run else 0,
_RUN_STALE_TIMEOUT,
)
def cancel_run(run_id: str) -> bool: def cancel_run(run_id: str) -> bool:
"""请求取消运行。""" """请求取消运行。"""
with _runs_lock: with _runs_lock:
...@@ -99,6 +134,8 @@ def get_running_for_target(target_id: str) -> Optional[dict]: ...@@ -99,6 +134,8 @@ def get_running_for_target(target_id: str) -> Optional[dict]:
用于防重复触发:同一目标已有未完成的巡检时返回该运行记录, 用于防重复触发:同一目标已有未完成的巡检时返回该运行记录,
前端据此决定是否阻止启动新巡检。 前端据此决定是否阻止启动新巡检。
""" """
# 兜底清理僵尸巡检,避免卡死线程遗留的未完成记录永远阻塞该目标
_cleanup_stale_runs()
with _runs_lock: with _runs_lock:
for run_id, run in _runs.items(): for run_id, run in _runs.items():
if (run.get("target_id") == target_id if (run.get("target_id") == target_id
...@@ -340,7 +377,9 @@ def run_inspection_sync(target_id: str, suite: str) -> dict: ...@@ -340,7 +377,9 @@ def run_inspection_sync(target_id: str, suite: str) -> dict:
try: try:
from . import notification_service from . import notification_service
report = report_service.get_report(report_id) report = report_service.get_report(report_id)
if report: if not report:
logger.warning("巡检完成通知:未找到报告 %s,跳过", report_id)
else:
# 构建报告链接(需要从配置获取外部访问地址) # 构建报告链接(需要从配置获取外部访问地址)
import os import os
host = os.environ.get('EXTERNAL_HOST', 'http://192.168.5.60:8088') host = os.environ.get('EXTERNAL_HOST', 'http://192.168.5.60:8088')
...@@ -348,7 +387,10 @@ def run_inspection_sync(target_id: str, suite: str) -> dict: ...@@ -348,7 +387,10 @@ def run_inspection_sync(target_id: str, suite: str) -> dict:
# 拼接访问 token,实现免登查看 # 拼接访问 token,实现免登查看
if report.get("access_token"): if report.get("access_token"):
report_url = f"{report_url}?token={report['access_token']}" report_url = f"{report_url}?token={report['access_token']}"
notification_service.send_notification(report, report_url) sent = notification_service.send_notification(report, report_url)
logger.info("巡检完成通知:%s(report_id=%s)",
"已发送" if sent else "未发送(无启用渠道或触发条件不满足)",
report_id)
# 连续异常告警检查 # 连续异常告警检查
try: try:
......
...@@ -12,7 +12,6 @@ statistics_service.py — 监测统计聚合 ...@@ -12,7 +12,6 @@ statistics_service.py — 监测统计聚合
from __future__ import annotations from __future__ import annotations
import logging
from collections import defaultdict from collections import defaultdict
from datetime import datetime, timedelta from datetime import datetime, timedelta
from typing import Dict, List, Optional, Tuple from typing import Dict, List, Optional, Tuple
...@@ -20,8 +19,9 @@ from typing import Dict, List, Optional, Tuple ...@@ -20,8 +19,9 @@ from typing import Dict, List, Optional, Tuple
from . import report_service, target_service from . import report_service, target_service
from ..utils.thresholds import _extract_numeric, THRESHOLDS from ..utils.thresholds import _extract_numeric, THRESHOLDS
from ..utils.display_names import DISPLAY_NAMES from ..utils.display_names import DISPLAY_NAMES
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.statistics") logger = get_logger("service_monitor.statistics")
NORMAL, WARNING, CRITICAL = "正常", "警告", "严重" NORMAL, WARNING, CRITICAL = "正常", "警告", "严重"
DEFAULT_RANGE_DAYS = 30 DEFAULT_RANGE_DAYS = 30
......
...@@ -9,7 +9,6 @@ target_service.py — 监测目标管理(CRUD + 连通性测试) ...@@ -9,7 +9,6 @@ target_service.py — 监测目标管理(CRUD + 连通性测试)
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import re import re
import uuid import uuid
from datetime import datetime from datetime import datetime
...@@ -18,8 +17,9 @@ from typing import Optional ...@@ -18,8 +17,9 @@ from typing import Optional
from ..utils.paths import TARGETS_FILE, ensure_dirs from ..utils.paths import TARGETS_FILE, ensure_dirs
from ..utils.crypto import encrypt_password, decrypt_password, is_encrypted from ..utils.crypto import encrypt_password, decrypt_password, is_encrypted
from ..utils.executor import LocalExecutor, SSHExecutor, SSHConnectionError, DependencyError from ..utils.executor import LocalExecutor, SSHExecutor, SSHConnectionError, DependencyError
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.target_service") logger = get_logger("service_monitor.target_service")
# 本地目标内置 ID # 本地目标内置 ID
LOCAL_TARGET_ID = "local" LOCAL_TARGET_ID = "local"
......
...@@ -19,11 +19,26 @@ if str(WEB_DIR) not in sys.path: ...@@ -19,11 +19,26 @@ if str(WEB_DIR) not in sys.path:
import pytest # noqa: E402 import pytest # noqa: E402
@pytest.fixture(autouse=True)
def _force_flask_page_fallback(monkeypatch):
"""页面路由强制走 Flask 模板回退,隔离本机 Vue dist 构建产物。
test_routes_sm.TestPages 的断言针对 Flask 回退渲染(含逐页鉴权逻辑:
未登录重定向、非管理员重定向、管理员可达)。而 server.py 在
frontend/dist 存在时会让 spa_fallback 直接返回 SPA index.html(对所有
请求一律 200,无服务端鉴权),使这些断言在本机构建了 Vue 时全部失败。
本仓库 frontend/dist 是 gitignore 的构建产物,不应影响测试结果,故统一
指向不存在的目录,让页面路由落到带鉴权的 Flask 回退分支。
"""
import server
monkeypatch.setattr(server, "VUE_DIST_DIR", Path("__nonexistent_vue_dist__"))
@pytest.fixture @pytest.fixture
def tmp_data(tmp_path, monkeypatch): def tmp_data(tmp_path, monkeypatch):
"""把 service_monitor 的数据目录重定向到 tmp_path,隔离测试。""" """把 service_monitor 的数据目录重定向到 tmp_path,隔离测试。"""
from service_monitor.utils import paths as sm_paths from service_monitor.utils import paths as sm_paths
from service_monitor.services import target_service, report_service from service_monitor.services import target_service, report_service, schedule_service
data_dir = tmp_path / "data" data_dir = tmp_path / "data"
reports_dir = data_dir / "reports" reports_dir = data_dir / "reports"
...@@ -31,15 +46,18 @@ def tmp_data(tmp_path, monkeypatch): ...@@ -31,15 +46,18 @@ def tmp_data(tmp_path, monkeypatch):
targets_file = data_dir / "targets.json" targets_file = data_dir / "targets.json"
report_index_file = data_dir / "report_index.json" report_index_file = data_dir / "report_index.json"
cooldown_file = data_dir / "alert_cooldown.json" cooldown_file = data_dir / "alert_cooldown.json"
schedules_file = data_dir / "schedules.json"
monkeypatch.setattr(sm_paths, "DATA_DIR", data_dir) monkeypatch.setattr(sm_paths, "DATA_DIR", data_dir)
monkeypatch.setattr(sm_paths, "REPORTS_DIR", reports_dir) monkeypatch.setattr(sm_paths, "REPORTS_DIR", reports_dir)
monkeypatch.setattr(sm_paths, "TARGETS_FILE", targets_file) monkeypatch.setattr(sm_paths, "TARGETS_FILE", targets_file)
monkeypatch.setattr(sm_paths, "REPORT_INDEX_FILE", report_index_file) monkeypatch.setattr(sm_paths, "REPORT_INDEX_FILE", report_index_file)
monkeypatch.setattr(sm_paths, "COOLDOWN_FILE", cooldown_file) monkeypatch.setattr(sm_paths, "COOLDOWN_FILE", cooldown_file)
monkeypatch.setattr(sm_paths, "SCHEDULES_FILE", schedules_file)
# service 模块在导入时已绑定常量引用,需同步 patch # service 模块在导入时已绑定常量引用,需同步 patch
monkeypatch.setattr(target_service, "TARGETS_FILE", targets_file) monkeypatch.setattr(target_service, "TARGETS_FILE", targets_file)
monkeypatch.setattr(report_service, "REPORTS_DIR", reports_dir) monkeypatch.setattr(report_service, "REPORTS_DIR", reports_dir)
monkeypatch.setattr(report_service, "REPORT_INDEX_FILE", report_index_file) monkeypatch.setattr(report_service, "REPORT_INDEX_FILE", report_index_file)
monkeypatch.setattr(report_service, "COOLDOWN_FILE", cooldown_file) monkeypatch.setattr(report_service, "COOLDOWN_FILE", cooldown_file)
monkeypatch.setattr(schedule_service, "SCHEDULES_FILE", schedules_file)
return {"data": data_dir, "reports": reports_dir, "targets": targets_file} return {"data": data_dir, "reports": reports_dir, "targets": targets_file}
...@@ -62,3 +62,71 @@ def test_list_asset_files(): ...@@ -62,3 +62,71 @@ def test_list_asset_files():
assert "common.sh" in names assert "common.sh" in names
assert "config.sh.template" in names assert "config.sh.template" in names
assert "01_system_basic.sh" in names assert "01_system_basic.sh" in names
# ============================================================
# 61_platform_log_check — 平台日志分析模块
# ============================================================
def test_platform_log_module_registered():
"""61 模块注册:仅 full 套件,quick 不含;asset 文件存在。"""
full = [m.id for m in cm.get_suite("full")]
quick = [m.id for m in cm.get_suite("quick")]
assert "61_platform_log_check" in full
assert "61_platform_log_check" not in quick
m = cm.get_module("61_platform_log_check")
assert m is not None
assert m.category == "system"
assert m.asset_path.name == "61_platform_log_check.sh"
assert m.asset_path.exists(), "61_platform_log_check.sh 资产文件缺失"
def test_render_config_platform_log_threshold_override():
"""目标级阈值覆盖 5 个新键追加在尾部。"""
target = {"thresholds": {
"platform_log_stale_warn_min": 20,
"platform_log_stale_crit_min": 90,
"platform_log_err_cnt_warn": 60,
"platform_log_err_cnt_crit": 120,
"platform_log_scan_lines": 3000,
}}
out = cm.render_config(target)
assert "export PLATFORM_LOG_STALE_WARN_MIN=20" in out
assert "export PLATFORM_LOG_STALE_CRIT_MIN=90" in out
assert "export PLATFORM_LOG_ERR_CNT_WARN=60" in out
assert "export PLATFORM_LOG_ERR_CNT_CRIT=120" in out
assert "export PLATFORM_LOG_SCAN_LINES=3000" in out
def test_render_config_platform_log_invalid_key_ignored():
"""非法阈值键被忽略(服务清单不开放目标级覆盖)。"""
target = {"thresholds": {
"platform_log_services_core": "evil:path", # 服务清单不允许覆盖
"platform_log_nonexistent": "x",
"cpu_warning": 91, # 合法键仍生效
}}
out = cm.render_config(target)
assert "PLATFORM_LOG_SERVICES_CORE=evil:path" not in out
assert "PLATFORM_LOG_NONEXISTENT" not in out
assert "export CPU_WARNING=91" in out
def test_platform_log_script_syntax():
"""61 脚本 bash -n 语法检查;环境无 bash 则降级为文件存在断言。"""
m = cm.get_module("61_platform_log_check")
assert m is not None
script = m.asset_path
assert script.exists()
import shutil
if shutil.which("bash"):
import subprocess
proc = subprocess.run(["bash", "-n", str(script)],
capture_output=True, text=True)
assert proc.returncode == 0, f"bash -n 语法错误:\n{proc.stderr}"
else:
# 无 bash 环境(如 Windows 裸机)降级:仅断言脚本存在且含关键输出函数
text = script.read_text(encoding="utf-8")
assert "output_result" in text
assert "detect_platform_type" in text
...@@ -11,15 +11,15 @@ check_modules.py — 检测模块清单与配置渲染 ...@@ -11,15 +11,15 @@ check_modules.py — 检测模块清单与配置渲染
from __future__ import annotations from __future__ import annotations
import logging
import shlex import shlex
from dataclasses import dataclass from dataclasses import dataclass
from pathlib import Path from pathlib import Path
from typing import Dict, List, Optional from typing import Dict, List, Optional
from .paths import CONFIG_TEMPLATE, COMMON_SH, SYSTEM_ASSETS_DIR, SERVICE_ASSETS_DIR from .paths import CONFIG_TEMPLATE, COMMON_SH, SYSTEM_ASSETS_DIR, SERVICE_ASSETS_DIR
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.check_modules") logger = get_logger("service_monitor.check_modules")
# 套件常量 # 套件常量
SUITE_QUICK = "quick" SUITE_QUICK = "quick"
...@@ -95,6 +95,8 @@ ALL_MODULES: List[CheckModule] = [ ...@@ -95,6 +95,8 @@ ALL_MODULES: List[CheckModule] = [
{SUITE_FULL}, "59_log_export_check.sh"), {SUITE_FULL}, "59_log_export_check.sh"),
CheckModule("60_repair_capability_check", "修复能力检测", "system", CheckModule("60_repair_capability_check", "修复能力检测", "system",
{SUITE_FULL}, "60_repair_capability_check.sh"), {SUITE_FULL}, "60_repair_capability_check.sh"),
CheckModule("61_platform_log_check", "平台日志分析", "system",
{SUITE_FULL}, "61_platform_log_check.sh"),
# ---- service 类 ---- # ---- service 类 ----
CheckModule("20_docker_basic", "Docker基础", "service", CheckModule("20_docker_basic", "Docker基础", "service",
{SUITE_QUICK, SUITE_FULL}, "20_docker_basic.sh"), {SUITE_QUICK, SUITE_FULL}, "20_docker_basic.sh"),
...@@ -236,6 +238,9 @@ def _is_valid_threshold_key(key: str) -> bool: ...@@ -236,6 +238,9 @@ def _is_valid_threshold_key(key: str) -> bool:
"disk_warning", "disk_critical", "disk_warning", "disk_critical",
"thread_warning", "thread_critical", "thread_warning", "thread_critical",
"java_log_path", "python_log_path", "nginx_log_path", "nacos_log_path", "java_log_path", "python_log_path", "nginx_log_path", "nacos_log_path",
"platform_log_stale_warn_min", "platform_log_stale_crit_min",
"platform_log_err_cnt_warn", "platform_log_err_cnt_crit",
"platform_log_scan_lines",
} }
return key.lower() in allowed return key.lower() in allowed
......
...@@ -15,7 +15,6 @@ executor.py — 检测执行器 ...@@ -15,7 +15,6 @@ executor.py — 检测执行器
from __future__ import annotations from __future__ import annotations
import abc import abc
import logging
import os import os
import socket import socket
import subprocess import subprocess
...@@ -30,8 +29,9 @@ from .paths import ( ...@@ -30,8 +29,9 @@ from .paths import (
COMMON_SH, CONFIG_TEMPLATE, COMMON_SH, CONFIG_TEMPLATE,
) )
from .check_modules import CheckModule from .check_modules import CheckModule
from utils.logger import get_logger
logger = logging.getLogger("service_monitor.executor") logger = get_logger("service_monitor.executor")
# 目标机上的工作目录前缀(后端根据 run_id 创建子目录) # 目标机上的工作目录前缀(后端根据 run_id 创建子目录)
_REMOTE_BASE = "/tmp/check_modules" _REMOTE_BASE = "/tmp/check_modules"
...@@ -313,6 +313,20 @@ class BaseExecutor(abc.ABC): ...@@ -313,6 +313,20 @@ class BaseExecutor(abc.ABC):
raise classify_ssh_error(e) raise classify_ssh_error(e)
def _copy_lf(src: Path, dst: Path) -> None:
"""复制文本资产文件并强制转 LF 换行。
Windows 开发机上 Git autocrlf=true 会让仓库工作区的 .sh 带 CRLF,
直接 shutil.copy2 会原样复制,本地 bash 执行时报 $'\r' 语法错误。
Git 提交与远端部署均要求 Linux 风格换行。
"""
try:
data = src.read_bytes()
dst.write_bytes(data.replace(b"\r\n", b"\n"))
except OSError as e:
raise IOError(f"复制资产文件失败 {src.name}: {e}") from e
# ============================================================ # ============================================================
# 本地执行器 # 本地执行器
# ============================================================ # ============================================================
...@@ -400,14 +414,14 @@ class LocalExecutor(BaseExecutor): ...@@ -400,14 +414,14 @@ class LocalExecutor(BaseExecutor):
lib_dir = self._local_workdir / "lib" lib_dir = self._local_workdir / "lib"
(lib_dir / "system").mkdir(parents=True, exist_ok=True) (lib_dir / "system").mkdir(parents=True, exist_ok=True)
(lib_dir / "service").mkdir(parents=True, exist_ok=True) (lib_dir / "service").mkdir(parents=True, exist_ok=True)
# 复制资产文件 # 复制资产文件(强制 LF,避免 Windows CRLF 在本地 bash 执行时报 $'\r' 语法错误)
shutil.copy2(COMMON_SH, lib_dir / "common.sh") _copy_lf(COMMON_SH, lib_dir / "common.sh")
for f in SYSTEM_ASSETS_DIR.iterdir(): for f in SYSTEM_ASSETS_DIR.iterdir():
if f.suffix == ".sh": if f.suffix == ".sh":
shutil.copy2(f, lib_dir / "system" / f.name) _copy_lf(f, lib_dir / "system" / f.name)
for f in SERVICE_ASSETS_DIR.iterdir(): for f in SERVICE_ASSETS_DIR.iterdir():
if f.suffix == ".sh": if f.suffix == ".sh":
shutil.copy2(f, lib_dir / "service" / f.name) _copy_lf(f, lib_dir / "service" / f.name)
self._workdir = str(self._local_workdir) self._workdir = str(self._local_workdir)
return self._workdir return self._workdir
...@@ -426,9 +440,11 @@ class LocalExecutor(BaseExecutor): ...@@ -426,9 +440,11 @@ class LocalExecutor(BaseExecutor):
# SSH 模式:通过 SFTP 写入 # SSH 模式:通过 SFTP 写入
self._get_ssh_executor()._write_config(wd, config_text) self._get_ssh_executor()._write_config(wd, config_text)
else: else:
# 本地模式:直接写入文件 # 本地模式:直接写入文件(统一 LF,避免 CRLF 破坏 bash 解析)
config_path = Path(wd) / "lib" / "config.sh" config_path = Path(wd) / "lib" / "config.sh"
config_path.write_text(config_text, encoding="utf-8") config_path.write_text(
config_text.replace("\r\n", "\n"), encoding="utf-8"
)
def _exec(self, command: str, timeout: int) -> str: def _exec(self, command: str, timeout: int) -> str:
if self._use_ssh: if self._use_ssh:
...@@ -654,7 +670,13 @@ class SSHExecutor(BaseExecutor): ...@@ -654,7 +670,13 @@ class SSHExecutor(BaseExecutor):
return self._sftp_client return self._sftp_client
def _exec(self, command: str, timeout: int) -> str: def _exec(self, command: str, timeout: int) -> str:
"""执行命令,如果 use_sudo=True 则通过 sudo 提权。""" """执行命令,如果 use_sudo=True 则通过 sudo 提权。
超时保护:paramiko 的 exec_command(timeout) 只作用于 socket I/O,
recv_exit_status() 会无限期等待远程进程退出。这里用
status_event.wait() 限制等待时间,超时则强制关闭 channel,
防止巡检线程被挂死的远程命令卡死(_runs 永不释放 → 后续 409)。
"""
# sudo 提权:用 bash -c 包裹命令 # sudo 提权:用 bash -c 包裹命令
if self.use_sudo: if self.use_sudo:
# 单引号内的命令不能再有单引号,需转义 # 单引号内的命令不能再有单引号,需转义
...@@ -662,9 +684,27 @@ class SSHExecutor(BaseExecutor): ...@@ -662,9 +684,27 @@ class SSHExecutor(BaseExecutor):
command = f"sudo bash -c '{escaped}'" command = f"sudo bash -c '{escaped}'"
ssh = self._get_ssh() ssh = self._get_ssh()
channel = None
try: try:
_, stdout, stderr = ssh.exec_command(command, timeout=timeout) _, stdout, stderr = ssh.exec_command(command, timeout=timeout)
exit_code = stdout.channel.recv_exit_status() channel = stdout.channel
# 等待远程进程退出,带超时(recv_exit_status 会无限阻塞)
if not channel.exit_status_ready():
channel.status_event.wait(timeout=timeout)
if not channel.exit_status_ready():
# 超时仍未退出,强制关闭 channel(远程命令随之断开)
logger.warning(
"SSH 命令超时 %ds,强制关闭 channel: %s",
timeout, command[:80])
try:
channel.close()
except Exception:
pass
return ""
exit_code = channel.recv_exit_status()
out = stdout.read().decode("utf-8", errors="replace") out = stdout.read().decode("utf-8", errors="replace")
err = stderr.read().decode("utf-8", errors="replace") err = stderr.read().decode("utf-8", errors="replace")
if err: if err:
...@@ -675,6 +715,12 @@ class SSHExecutor(BaseExecutor): ...@@ -675,6 +715,12 @@ class SSHExecutor(BaseExecutor):
return out return out
except Exception as e: except Exception as e:
logger.warning("SSH 执行失败: %s (cmd: %s)", e, command[:60]) logger.warning("SSH 执行失败: %s (cmd: %s)", e, command[:60])
# 兜底关闭 channel,避免连接泄漏
if channel is not None:
try:
channel.close()
except Exception:
pass
return "" return ""
def _upload_file(self, local_path: Path, remote_path: str): def _upload_file(self, local_path: Path, remote_path: str):
...@@ -689,17 +735,22 @@ class SSHExecutor(BaseExecutor): ...@@ -689,17 +735,22 @@ class SSHExecutor(BaseExecutor):
ssh = self._get_ssh() ssh = self._get_ssh()
_, stdout, _ = ssh.exec_command(f"mkdir -p {remote_dir}", timeout=10) _, stdout, _ = ssh.exec_command(f"mkdir -p {remote_dir}", timeout=10)
stdout.channel.recv_exit_status() stdout.channel.recv_exit_status()
sftp.put(str(local_path), remote_path) # 关键:以二进制读入并强制转 LF 换行再上传
# (Windows 开发机上 Git autocrlf 会让 .sh 源码带 CRLF,sftp.put
# 会原样上传,远端 bash 解析报 $'\r' 语法错误,模块检测项全部丢失)
data = local_path.read_bytes().replace(b"\r\n", b"\n")
with sftp.open(remote_path, "wb") as f:
f.write(data)
logger.debug("上传: %s → %s", local_path.name, remote_path) logger.debug("上传: %s → %s", local_path.name, remote_path)
except Exception as e: except Exception as e:
raise IOError(f"上传文件失败 {local_path.name}: {e}") from e raise IOError(f"上传文件失败 {local_path.name}: {e}") from e
def _write_config(self, wd: str, config_text: str) -> None: def _write_config(self, wd: str, config_text: str) -> None:
# SSH 模式:通过 SFTP 写入渲染后的 config.sh # SSH 模式:通过 SFTP 写入渲染后的 config.sh(强制 LF)
sftp = self._get_sftp() sftp = self._get_sftp()
remote_config = f"{wd}/lib/config.sh" remote_config = f"{wd}/lib/config.sh"
with sftp.open(remote_config, "w") as f: with sftp.open(remote_config, "w") as f:
f.write(config_text) f.write(config_text.replace("\r\n", "\n"))
def cleanup(self) -> None: def cleanup(self) -> None:
super().cleanup() super().cleanup()
......
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论