Update to 7.7.0

Since version 7.7.0, we recommend yours update python to 3.12.

[+] Using nginx technology to load static files improves access speed
[+] Refactor homepage, website, FTP, and database using vue3
[+] Table loading changed to skeleton screen
[+] Add Website statistics-v2 professional plug-in
[+] Add Home page - top 5 resource occupancy
[+] Add protection for Files management (requires Tamper-proof for Enterprise 3.7)
[+] Website, FTP, Databases page add program status
[+] Add FTP log analysis (only supports Centos)
[+] Add password-free login to phpMyAdmin
[+] Add Proxy Project in Website (Supported when web service uses Nginx)
[+] Add WP Toolkit (Pro version only)
[+] Redesigned Docker module
[+] Add WP Toolkit Protection
[+] Add WP Toolkit Backup and Restore
[+] Add WP Toolkit Migrated
[+] Add WP Toolkit Clone site (supports new domain and subdomain)
[+] Add WP Toolkit Create site from backup of other panel
[+] Add WP Toolkit support for Cron automatic backup (only save Local disk)
[+] Add WP Toolkit operation log
[+] Add Integrity check for WP Toolkit
[+] Add WP Toolkit plug-in management and themes management

[*] Optimize phpMyAdmin formula access method
[*] Optimize Home page PHP display problem
[*] Optimize jump to the login interface after the login expires
[*] Optimize automatic renewal of SSL at some times
[*] Optimize Let's Encrypt to increase application success rate

[-] Fix Logs Audit cannot be opened
[-] Fix apache URL rewrite issue
[-] Fix phpmyadmin installation problem
[-] Fix the problem that some servers cannot install software
[-] Fix upload file error
[-] Fix left menu hiding problem
[-] Fix aaPanel Mobile QR code display problem
[-] Fix problem that third-party plug-ins are not displayed in the App Store
[-] Fix issue where the menu bar is blank when opening new tabs
[-] Fixed panel not being accessible in some cases
[-] Fix the issue where Curl warning caused the inability to apply for SSL
[-] Fix Quota issues for Website, FTP, Databases
[-] Fix file interface display problem on mobile terminal
This commit is contained in:
Jack
2024-07-19 11:25:10 +08:00
parent ed34994f45
commit ed55fa708d
1949 changed files with 329210 additions and 26644 deletions
+495
View File
@@ -0,0 +1,495 @@
import json
import os
from typing import Dict, Union
from .mods import TaskConfig, TaskTemplateConfig, TaskRecordConfig, SenderConfig, load_task_template_by_config, \
load_task_template_by_file, UPDATE_MOD_PUSH_FILE, UPDATE_VERSION_FILE, PUSH_DATA_PATH
from .base_task import BaseTask
from .send_tool import WxAccountMsg, WxAccountLoginMsg, WxAccountMsgBase
from .system import PushSystem, get_push_public_data, push_by_task_keyword, push_by_task_id
from .manager import PushManager
from .util import read_file, write_file
__all__ = [
"TaskConfig",
"TaskTemplateConfig",
"TaskRecordConfig",
"SenderConfig",
"load_task_template_by_config",
"load_task_template_by_file",
"BaseTask",
"WxAccountMsg",
"WxAccountLoginMsg",
"WxAccountMsgBase",
"PushSystem",
"get_push_public_data",
"PushManager",
"push_by_task_keyword",
"push_by_task_id",
"UPDATE_MOD_PUSH_FILE",
"update_mod_push_system",
"UPDATE_VERSION_FILE",
"PUSH_DATA_PATH",
"get_default_module_dict",
]
def update_mod_push_system():
if os.path.exists(UPDATE_MOD_PUSH_FILE):
return
# 只将已有的告警任务("site_push", "system_push", "database_push") 移动
try:
push_data = json.loads(read_file("/www/server/panel/class/push/push.json"))
except:
return
if not isinstance(push_data, dict):
return
pmgr = PushManager()
default_module_dict = get_default_module_dict()
for key, value in push_data.items():
if key == "site_push":
_update_site_push(value, pmgr, default_module_dict)
elif key == "system_push":
_update_system_push(value, pmgr, default_module_dict)
elif key == "database_push":
_update_database_push(value, pmgr, default_module_dict)
elif key == "rsync_push":
_update_rsync_push(value, pmgr, default_module_dict)
elif key == "load_balance_push":
_update_load_push(value, pmgr, default_module_dict)
elif key == "task_manager_push":
_update_task_manager_push(value, pmgr, default_module_dict)
write_file(UPDATE_MOD_PUSH_FILE, "")
def get_default_module_dict():
res = {}
wx_account_list = []
for data in SenderConfig().config:
if not data["used"]:
continue
if data.get("original", False):
res[data["sender_type"]] = data["id"]
if data["sender_type"] == "webhook":
res[data["data"].get("title")] = data["id"]
if data["sender_type"] == "wx_account":
wx_account_list.append(data)
wx_account_list.sort(key=lambda x: x.get("data", {}).get("create_time", ""))
if wx_account_list:
res["wx_account"] = wx_account_list[0]["id"]
return res
def _update_site_push(old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
if v["type"] == "ssl":
push_data = {
"template_id": "1",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", "all"),
"cycle": v.get("cycle", 15)
},
"number_rule": {
"total": v.get("push_count", 1)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "site_endtime":
push_data = {
"template_id": "2",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 7)
},
"number_rule": {
"total": v.get("push_count", 1)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "panel_pwd_endtime":
push_data = {
"template_id": "3",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 15),
"interval": 600
},
"number_rule": {
"total": v.get("push_count", 1)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "ssh_login_error":
push_data = {
"template_id": "4",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 30),
"count": v.get("count", 3),
"interval": v.get("interval", 600)
},
"number_rule": {
"day_num": v.get("day_limit", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "services":
push_data = {
"template_id": "5",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", "nginx"),
"count": v.get("count", 3),
"interval": v.get("interval", 600)
},
"number_rule": {
"day_num": v.get("day_limit", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "panel_safe_push":
push_data = {
"template_id": "6",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {},
"number_rule": {
"day_num": v.get("day_limit", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "ssh_login":
push_data = {
"template_id": "7",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "panel_login":
push_data = {
"template_id": "8",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "project_status":
push_data = {
"template_id": "9",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 1),
"project": v.get("project", 0),
"count": v.get("count", 2) if v.get("count", 2) not in (1, 2) else 2,
"interval": v.get("interval", 600)
},
"number_rule": {
"day_num": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "panel_update":
push_data = {
"template_id": "10",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {},
"number_rule": {
"day_num": 1
}
}
}
pmgr.set_task_conf_data(push_data)
send_type = None
login_send_type_conf = "/www/server/panel/data/panel_login_send.pl"
if os.path.exists(login_send_type_conf):
send_type = read_file(login_send_type_conf).strip()
else:
# 兼容之前的
if os.path.exists("/www/server/panel/data/login_send_type.pl"):
send_type = read_file("/www/server/panel/data/login_send_type.pl")
else:
if os.path.exists('/www/server/panel/data/login_send_mail.pl'):
send_type = "mail"
if os.path.exists('/www/server/panel/data/login_send_dingding.pl'):
send_type = "dingding"
if isinstance(send_type, str):
sender_list = [df_mdl[i.strip()] for i in send_type.split(",") if i.strip() in df_mdl]
push_data = {
"template_id": "8",
"task_data": {
"status": True,
"sender": sender_list,
"task_data": {},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
login_send_type_conf = "/www/server/panel/data/ssh_send_type.pl"
if os.path.exists(login_send_type_conf):
ssh_send_type = read_file(login_send_type_conf).strip()
if isinstance(ssh_send_type, str):
sender_list = [df_mdl[i.strip()] for i in ssh_send_type.split(",") if i.strip() in df_mdl]
push_data = {
"template_id": "7",
"task_data": {
"status": True,
"sender": sender_list,
"task_data": {},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
return
def _update_system_push(old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
if v["type"] == "disk":
push_data = {
"template_id": "20",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", "/"),
"cycle": v.get("cycle", 2) if v.get("cycle", 2) not in (1, 2) else 2,
"count": v.get("count", 80),
},
"number_rule": {
"total": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
if v["type"] == "disk":
push_data = {
"template_id": "21",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 5) if v.get("cycle", 5) not in (3, 5, 15) else 5,
"count": v.get("count", 80),
},
"number_rule": {
"total": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
if v["type"] == "load":
push_data = {
"template_id": "22",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 5) if v.get("cycle", 5) not in (1, 5, 15) else 5,
"count": v.get("count", 80),
},
"number_rule": {
"total": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
if v["type"] == "mem":
push_data = {
"template_id": "23",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"cycle": v.get("cycle", 5) if v.get("cycle", 5) not in (3, 5, 15) else 5,
"count": v.get("count", 80),
},
"number_rule": {
"total": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
return
def _update_database_push(old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
if v["type"] == "mysql_pwd_endtime":
push_data = {
"template_id": "30",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", []),
"cycle": v.get("cycle", 15),
},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
elif v["type"] == "mysql_replicate_status":
push_data = {
"template_id": "31",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", []),
"count": v.get("cycle", 15),
"interval": v.get("interval", 600)
},
"number_rule": {}
}
}
pmgr.set_task_conf_data(push_data)
return None
def _update_rsync_push(
old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
push_data = {
"template_id": "40",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"interval": v.get("interval", 600)
},
"number_rule": {
"day_num": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
def _update_load_push(
old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
push_data = {
"template_id": "50",
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", ""),
"cycle": v.get("cycle", "200|301|302|403|404")
},
"number_rule": {
"day_num": v.get("push_count", 2)
}
}
}
pmgr.set_task_conf_data(push_data)
def _update_task_manager_push(
old_data: Dict[str, Dict[str, Union[str, int, float, list]]],
pmgr: PushManager,
df_mdl: Dict[str, str]):
for k, v in old_data.items():
sender_list = [df_mdl[i.strip()] for i in v.get("module", "").split(",") if i.strip() in df_mdl]
template_id_dict = {
"task_manager_cpu": "60",
"task_manager_mem": "61",
"task_manager_process": "62"
}
if v["type"] in template_id_dict:
push_data = {
"template_id": template_id_dict[v["type"]],
"task_data": {
"status": bool(v.get("status", True)),
"sender": sender_list,
"task_data": {
"project": v.get("project", ""),
"count": v.get("count", 80),
"interval": v.get("count", 600),
},
"number_rule": {
"day_num": v.get("push_count", 3)
}
}
}
pmgr.set_task_conf_data(push_data)
+207
View File
@@ -0,0 +1,207 @@
from typing import Union, Optional, List, Tuple
from .send_tool import WxAccountMsg
# 告警系统在处理每个任务时,都会重新建立有一个Task的对象,(请勿在__init__的初始化函数中添加任何参数)
# 故每个对象中都可以大胆存放本任务所有数据,不会影响同类型的其他任务
class BaseTask:
def __init__(self):
self.source_name: str = ''
self.title: str = '' # 这个是告警任务的标题(根据实际情况改变)
self.template_name: str = '' # 这个告警模板的标题(不会改变)
def check_task_data(self, task_data: dict) -> Union[dict, str]:
"""
检查设置的告警参数(是否合理)
@param task_data: 传入的告警参数,提前会经过默认值处理(即没有的字段添加默认值)
@return: 当检查无误时,返回一个 dict 当做后续的添加和修改的数据,
当检查有误时, 直接返回错误信息的字符串
"""
raise NotImplementedError()
def get_keyword(self, task_data: dict) -> str:
"""
返回一个关键字,用于后续查询或执行任务时使用, 例如:防篡改告警,可以根据其规则id生成一个关键字,
后续通过规则id和来源tamper 查询并使用
@param task_data: 通过check_args后生成的告警参数字典
@return: 返回一个关键词字符串
"""
raise NotImplementedError()
def get_title(self, task_data: dict) -> str:
"""
返回一个标题
@param task_data: 通过check_args后生成的告警参数字典
@return: 返回一个关键词字符串
"""
if self.title:
return self.title
return self.template_name
def task_run_end_hook(self, res: dict) -> None:
"""
在告警系统中。执行完了任务后,会去掉用这个函数
@type res: dict, 执行任务的结果
@return:
"""
return
def task_config_update_hook(self, task: dict) -> None:
"""
在告警管理中。更新任务数据后,会去掉用这个函数
@return:
"""
return
def task_config_remove_hook(self, task: dict) -> None:
"""
在告警管理中。移除这个任务后,会去掉用这个函数
@return:
"""
return
def task_config_create_hook(self, task: dict) -> None:
"""
在告警管理中。新建这个任务后,会去掉用这个函数
@return:
"""
return
def check_time_rule(self, time_rule: dict) -> Union[dict, str]:
"""
检查和修改设置的告警的时间控制参数是是否合理
可以添加参数 get_by_func 字段用于指定使用本类中的那个函数执行时间判断标准, 替换标准的时间规则判断功能
↑示例如本类中的: can_send_by_time_rule
@param time_rule: 传入的告警参数,提前会经过默认值处理(即没有的字段添加默认值)
@return: 当检查无误时,返回一个 dict 当做后续的添加和修改的数据,
当检查有误时, 直接返回错误信息的字符串
"""
return time_rule
def check_num_rule(self, num_rule: dict) -> Union[dict, str]:
"""
检查和修改设置的告警的次数控制参数是是否合理
可以添加参数 get_by_func 字段用于指定使用本类中的那个函数执行次数判断标准, 替换标准的次数规则判断功能
↑示例如本类中的: can_send_by_num_rule
@param num_rule: 传入的告警参数,提前会经过默认值处理(即没有的字段添加默认值)
@return: 当检查无误时,返回一个 dict 当做后续的添加和修改的数据,
当检查有误时, 直接返回错误信息的字符串
"""
return num_rule
def can_send_by_num_rule(self, task_id: str, task_data: dict, number_rule: dict, push_data: dict) -> Optional[str]:
"""
这是一个通过函数判断是否能够发送告警的示例,并非每一个告警任务都需要有
@param task_id: 任务id
@param task_data: 告警参数信息
@param number_rule: 次数控制信息
@param push_data: 本次要发送的告警信息的原文,应当为字典, 来自 get_push_data 函数的返回值
@return: 返回None
"""
return None
def can_send_by_time_rule(self, task_id: str, task_data: dict, time_rule: dict, push_data: dict) -> Optional[str]:
"""
这是一个通过函数判断是否能够发送告警的示例,并非每一个告警任务都需要有
@param task_id: 任务id
@param task_data: 告警参数信息
@param time_rule: 时间控制信息
@param push_data: 本次要发送的告警信息的原文,应当为字典, 来自 get_push_data 函数的返回值
@return:
"""
return None
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
"""
判断这个任务是否需要返送
@param task_id: 任务id
@param task_data: 任务的告警参数
@return: 如果触发了告警,返回一个dict的原文,作为告警信息,否则应当返回None表示未触发
返回之中应当包含一个 msg_list 的键(值为List[str]类型),将主要的信息返回
用于以下信息的自动序列化包含[dingding, feishu, mail, weixin, web_hook]
短信和微信公众号由于长度问题,必须每个任务手动实现
"""
raise NotImplementedError()
def filter_template(self, template: dict) -> Optional[dict]:
"""
过滤 和 更改模板中的信息, 返回空表是当前无法设置该任务
@param template: 任务的模板信息
@return:
"""
raise NotImplementedError()
# push_public_data 公共的告警参数提取位置
# 内容包含:
# ip 网络ip
# local_ip 本机ip
# time 时间日志的字符串
# timestamp 当前的时间戳
# server_name 服务器别名
def to_dingding_msg(self, push_data: dict, push_public_data: dict) -> str:
print("dddddddddddddddddddddddddddddddddddddddddd")
msg_list = push_data.get('msg_list', None)
if msg_list is None:
raise ValueError("任务:{}的告警推送数据参数错误, 没有msg_list字段".format(self.title))
print("dddddddddddddddddddddddddddddddddddddddddd")
return self.public_headers_msg(push_public_data,dingding=True) + "\n\n" + "\n\n".join(msg_list)
def to_feishu_msg(self, push_data: dict, push_public_data: dict) -> str:
msg_list = push_data.get('msg_list', None)
if msg_list is None:
raise ValueError("任务:{}的告警推送数据参数错误, 没有msg_list字段".format(self.title))
return self.public_headers_msg(push_public_data) + "\n\n" + "\n\n".join(msg_list)
def to_mail_msg(self, push_data: dict, push_public_data: dict) -> str:
msg_list = push_data.get('msg_list', None)
if msg_list is None:
raise ValueError("任务:{}的告警推送数据参数错误, 没有msg_list字段".format(self.title))
public_headers = self.public_headers_msg(push_public_data, "<br>")
return public_headers + "<br>" + "<br>".join(msg_list)
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
"""
返回 短信告警的类型和数据
@param push_data:
@param push_public_data:
@return: 第一项是类型, 第二项是数据
"""
raise NotImplementedError()
def to_weixin_msg(self, push_data: dict, push_public_data: dict) -> str:
msg_list = push_data.get('msg_list', None)
if msg_list is None:
raise ValueError("任务:{}的告警推送数据参数错误, 没有msg_list字段".format(self.title))
spc = "\n "
public_headers = self.public_headers_msg(push_public_data, "\n ")
return public_headers + spc + spc.join(msg_list)
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
raise NotImplementedError()
def to_web_hook_msg(self, push_data: dict, push_public_data: dict) -> str:
msg_list = push_data.get('msg_list', None)
if msg_list is None:
raise ValueError("任务:{}的告警推送数据参数错误, 没有msg_list字段".format(self.title))
public_headers = self.public_headers_msg(push_public_data, "\n")
return public_headers + "\n" + "\n".join(msg_list)
def public_headers_msg(self, push_public_data: dict, spc: str = None,dingding=False) -> str:
if spc is None:
spc = "\n\n"
title = self.title
print(title)
if dingding:
print("dingdingtitle",title)
if "面板" not in title:
title += "面板"
print("dingdingtitle",title)
print(title)
return spc.join([
"#### {}".format(title),
">服务器:" + push_public_data['server_name'],
">IP地址:{}(外) {}(内)".format(push_public_data['ip'], push_public_data['local_ip']),
">发送时间:" + push_public_data['time']
])
+27
View File
@@ -0,0 +1,27 @@
import os
from .util import read_file, write_file
def rsync_compatible():
files = [
"/www/server/panel/class/push/rsync_push.py",
"/www/server/panel/plugin/rsync/rsync_push.py",
]
for f in files:
print(f)
if not os.path.exists(f):
continue
src_data = read_file(f)
if src_data.find("push_rsync_by_task_name") != -1:
continue
src_data = src_data.replace("""if __name__ == "__main__":
rsync_push().main()""", """
if __name__ == "__main__":
try:
sys.path.insert(0, "/www/server/panel")
from mod.base.push_mod.rsync_push import push_rsync_by_task_name
push_rsync_by_task_name(sys.argv[1])
except:
rsync_push().main()
""")
write_file(f, src_data)
+239
View File
@@ -0,0 +1,239 @@
import json
import os
import sys
import ipaddress
from datetime import datetime, timedelta
from typing import Tuple, Union, Optional
from .send_tool import WxAccountMsg
from .base_task import BaseTask
from .util import read_file, DB, GET_CLASS
try:
if "/www/server/panel/class" not in sys.path:
sys.path.insert(0, "/www/server/panel/class")
from panel_msg.collector import DatabasePushMsgCollect
except ImportError:
DatabasePushMsgCollect = None
def is_ipaddress(ip_data: str) -> bool:
try:
ipaddress.ip_address(ip_data)
except ValueError:
return False
return True
class MysqlPwdEndTimeTask(BaseTask):
def __init__(self):
super().__init__()
self.template_name = "MySQL数据库密码到期"
self.source_name = "mysql_pwd_end"
self.push_db_user = ""
def get_title(self, task_data: dict) -> str:
return "Msql:" + task_data["project"][1] + "用户密码到期提醒"
def check_task_data(self, task_data: dict) -> Union[dict, str]:
task_data["interval"] = 600
if not (isinstance(task_data["project"], list) and len(task_data["project"]) == 3):
return "设置的用户格式错误"
project = task_data["project"]
if not (isinstance(project[0], int) and isinstance(project[1], str) and is_ipaddress(project[2])):
return "设置的检测用户格式错误"
if not (isinstance(task_data["cycle"], int) and task_data["cycle"] >= 1):
return "到期时间参数错误,至少为 1 天"
return task_data
def get_keyword(self, task_data: dict) -> str:
return "_".join([str(i) for i in task_data["project"]])
def check_num_rule(self, num_rule: dict) -> Union[dict, str]:
num_rule["day_num"] = 1
return num_rule
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
sid = task_data["project"][0]
username = task_data["project"][1]
host = task_data["project"][2]
if "/www/server/panel/class" not in sys.path:
sys.path.insert(0, "/www/server/panel/class")
try:
import panelMysql
import db_mysql
except ImportError:
return None
if sid == 0:
try:
db_port = int(panelMysql.panelMysql().query("show global variables like 'port'")[0][1])
if db_port == 0:
db_port = 3306
except:
db_port = 3306
conn_config = {
"db_host": "localhost",
"db_port": db_port,
"db_user": "root",
"db_password": DB("config").where("id=?", (1,)).getField("mysql_root"),
"ps": "本地服务器",
}
else:
conn_config = DB("database_servers").where("id=? AND LOWER(db_type)=LOWER('mysql')", (sid,)).find()
if not conn_config:
return None
mysql_obj = db_mysql.panelMysql().set_host(conn_config["db_host"], conn_config["db_port"], None,
conn_config["db_user"], conn_config["db_password"])
if isinstance(mysql_obj, bool):
return None
data_list = mysql_obj.query(
"SELECT password_last_changed FROM mysql.user WHERE user='{}' AND host='{}';".format(username, host))
if not isinstance(data_list, list) or not data_list:
return None
try:
# todo:检查这里的时间转化逻辑问题
last_time = data_list[0][0]
expire_time = last_time + timedelta(days=task_data["cycle"])
except:
return None
if datetime.now() > expire_time:
self.title = self.get_title(task_data)
self.push_db_user = username
return {"msg_list": [
">告警类型:MySQL密码即将到期",
">告警内容:{} {}@{} 密码过期时间<font color=#ff0000>{} 天</font>".format(
conn_config["ps"], username, host, expire_time.strftime("%Y-%m-%d %H:%M:%S"))
]}
def filter_template(self, template: dict) -> Optional[dict]:
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return "", {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "MySQL数据库密码到期"
msg.msg = "Mysql用户:{}的密码即将过期,请注意".format(self.push_db_user)
msg.next_msg = "请登录面板,查看主机情况"
return msg
class MysqlReplicateStatusTask(BaseTask):
def __init__(self):
super().__init__()
self.template_name = "MySQL主从复制异常告警"
self.source_name = "mysql_replicate_status"
self.title = "MySQL主从复制异常告警"
self.slave_ip = ''
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not (isinstance(task_data["project"], str) and task_data["project"]):
return "请选择告警的从库!"
if not (isinstance(task_data["count"], int) and task_data["count"] in (1, 2)):
return "是否自动修复选择错误!"
if not (isinstance(task_data["interval"], int) and task_data["interval"] >= 60):
return "检查间隔时间错误,至少需要60s的间隔"
return task_data
def get_keyword(self, task_data: dict) -> str:
return task_data["project"]
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
import PluginLoader
args = GET_CLASS()
args.slave_ip = task_data["project"]
res = PluginLoader.plugin_run("mysql_replicate", "get_replicate_status", args)
if res.get("status", False) is False:
return None
self.slave_ip = task_data["project"]
if len(res.get("data", [])) == 0:
s_list = [">告警类型:MySQL主从复制异常告警",
">告警内容:<font color=#ff0000>从库 {} 主从复制已停止,请尽快登录面板查看详情</font>".format(
task_data["project"])]
return {"msg_list": s_list}
sql_status = io_status = False
for item in res.get("data", []):
if item["name"] == "Slave_IO_Running" and item["value"] == "Yes":
io_status = True
if item["name"] == "Slave_SQL_Running" and item["value"] == "Yes":
sql_status = True
if io_status is True and sql_status is True:
break
if io_status is False or sql_status is False:
repair_txt = "请尽快登录面板查看详情"
if task_data["count"] == 1: # 自动修复
PluginLoader.plugin_run("mysql_replicate", "repair_replicate", args)
repair_txt = ",正在尝试修复"
s_list = [">告警类型:MySQL主从复制异常告警",
">告警内容:<font color=#ff0000>从库 {} 主从复制发生异常{}</font>".format(
task_data["project"], repair_txt)]
return {"msg_list": s_list}
return None
@staticmethod
def _get_mysql_replicate():
slave_list = []
mysql_replicate_path = os.path.join("/www/server/panel/plugin", "mysql_replicate", "config.json")
if os.path.isfile(mysql_replicate_path):
conf = read_file(mysql_replicate_path)
try:
conf = json.loads(conf)
slave_list = [{"title": slave_ip, "value": slave_ip} for slave_ip in conf["slave"].keys()]
except:
pass
return slave_list
def filter_template(self, template: dict) -> Optional[dict]:
template["field"][0]["items"] = self._get_mysql_replicate()
if not template["field"][0]["items"]:
return None
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "MySQL主从复制异常告警"
msg.msg = "从库 {} 主从复制发生异常".format(self.slave_ip)
msg.next_msg = "请登录面板,在[软件商店-MySQL主从复制(重构版)]中查看"
return msg
class ViewMsgFormat(object):
_FORMAT = {
"30": (
lambda x: "<span>剩余时间小于{}天{}</span>".format(
x["task_data"].get("cycle"),
("(如未处理,次日会重新发送1次,持续%d天)" % x.get("number_rule", {}).get("day_num", 0)) if x.get(
"number_rule", {}).get("day_num", 0) else ""
)
),
"31": (lambda x: "<span>MySQL主从复制异常告警</span>".format()),
}
def get_msg(self, task: dict) -> Optional[str]:
if task["template_id"] in self._FORMAT:
return self._FORMAT[task["template_id"]](task)
return None
@@ -0,0 +1,154 @@
[
{
"id": "30",
"ver": "1",
"used": true,
"source": "mysql_pwd_end",
"title": "MySQL数据库密码到期",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.database_push",
"name": "MysqlPwdEndTimeTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "选择用户",
"type": "cascader",
"default": 0,
"items": [
{
"url": "database?action=GetPushUser"
},
{
"url": "database?action=GetPushUser",
"data": [
"sid"
]
},
{
"url": "database?action=GetPushUser",
"data": [
"sid",
"username"
]
}
]
},
{
"attr": "cycle",
"name": "剩余天数",
"type": "number",
"unit": "天",
"suffix": "",
"default": 15
}
],
"sorted": [
[
"project"
],
[
"cycle"
]
]
},
"default": {
"project": [],
"cycle": 15
},
"advanced_default": {
"number_rule": {
"total": 2
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": false
},
{
"id": "31",
"ver": "1",
"used": true,
"source": "mysql_replicate_status",
"title": "Msql主从同步告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.database_push",
"name": "MysqlReplicateStatusTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "选择监控的从库",
"type": "select",
"default": null,
"items": []
},
{
"attr": "count",
"name": "自动修复",
"type": "radio",
"suffix": "",
"default": 1,
"items": [
{
"title": "自动尝试修复",
"value": 1
},
{
"title": "不做修复尝试",
"value": 2
}
]
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"project"
],
[
"count"
],
[
"interval"
]
]
},
"default": {
"project": "",
"count": 2,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": false
}
]
+271
View File
@@ -0,0 +1,271 @@
import json
import os
import time
from typing import Tuple, Union, Optional
from .mods import PUSH_DATA_PATH, TaskTemplateConfig
from .send_tool import WxAccountMsg
from .base_task import BaseTask
from .util import read_file, DB, GET_CLASS, write_file
class NginxLoadTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "nginx_load_push"
self.template_name = "负载均衡告警"
self._tip_counter = None
@property
def tip_counter(self) -> dict:
if self._tip_counter is not None:
return self._tip_counter
tip_counter = '{}/load_balance_push.json'.format(PUSH_DATA_PATH)
if os.path.exists(tip_counter):
try:
self._tip_counter = json.loads(read_file(tip_counter))
except json.JSONDecodeError:
self._tip_counter = {}
else:
self._tip_counter = {}
return self._tip_counter
def save_tip_counter(self):
tip_counter = '{}/load_balance_push.json'.format(PUSH_DATA_PATH)
write_file(tip_counter, json.dumps(self.tip_counter))
def get_title(self, task_data: dict) -> str:
if task_data["project"] == "all":
return "负载节点异常告警"
return "负载节点[{}]异常告警".format(task_data["project"])
def check_task_data(self, task_data: dict) -> Union[dict, str]:
all_upstream_name = DB("upstream").field("name").select()
if isinstance(all_upstream_name, str) and all_upstream_name.startswith("error"):
return '没有负载均衡配置,无法设置告警'
all_upstream_name = [i["name"] for i in all_upstream_name]
if not bool(all_upstream_name):
return '没有负载均衡配置,无法设置告警'
if task_data["project"] not in all_upstream_name and task_data["project"] != "all":
return '没有该负载均衡配置,无法设置告警'
cycle = []
for i in task_data["cycle"].split("|"):
if bool(i) and i.isdecimal():
code = int(i)
if 100 <= code < 600:
cycle.append(str(code))
if not bool(cycle):
return '没有指定任何错误码,无法设置告警'
task_data["cycle"] = "|".join(cycle)
return task_data
def get_keyword(self, task_data: dict) -> str:
return task_data["project"]
def _check_func(self, upstream_name: str, codes: str) -> list:
import PluginLoader
get_obj = GET_CLASS()
get_obj.upstream_name = upstream_name
# 调用外部插件检查负载均衡的健康状况
upstreams = PluginLoader.plugin_run("load_balance", "get_check_upstream", get_obj)
access_codes = [int(i) for i in codes.split("|") if bool(i.strip())]
res_list = []
for upstream in upstreams:
# 检查每个节点,返回有问题的节点信息
res = upstream.check_nodes(access_codes, return_nodes=True)
for ping_url in res:
if ping_url in self.tip_counter:
self.tip_counter[ping_url].append(int(time.time()))
idx = 0
for i in self.tip_counter[ping_url]:
# 清理超过4分钟的记录
if time.time() - i > 60 * 4:
idx += 1
self.tip_counter[ping_url] = self.tip_counter[ping_url][idx:]
print("self.tip_counter[ping_url]",self.tip_counter[ping_url])
# 如果一个节点连续三次出现在告警列表中,则视为需要告警
if len(self.tip_counter[ping_url]) >= 3:
res_list.append(ping_url)
self.tip_counter[ping_url] = []
else:
self.tip_counter[ping_url] = [int(time.time()), ]
self.save_tip_counter()
return res_list
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
err_nodes = self._check_func(task_data["project"], task_data["cycle"])
if not err_nodes:
return None
pj = "负载均衡:【{}】".format(task_data["project"]) if task_data["project"] != "all" else "负载均衡"
nodes = '、'.join(err_nodes)
return {
"msg_list": [
">通知类型:企业版负载均衡告警",
">告警内容:<font color=#ff0000>{}配置下的节点【{}】出现访问错误,请及时关注节点情况并处理。</font> ".format(
pj, nodes),
],
"pj": pj,
"nodes": nodes
}
def filter_template(self, template: dict) -> Optional[dict]:
if not os.path.exists("/www/server/panel/plugin/load_balance/load_balance_main.py"):
return None
all_upstream = DB("upstream").field("name").select()
if isinstance(all_upstream, str) and all_upstream.startswith("error"):
return None
all_upstream_name = [i["name"] for i in all_upstream]
if not all_upstream_name:
return None
for name in all_upstream_name:
template["field"][0]["items"].append({
"title": name,
"value": name
})
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "负载均衡告警"
msg.msg = "负载均衡出现节点异常,请登录面板查看"
return msg
def task_config_create_hook(self, task: dict) -> None:
old_config_file = "/www/server/panel/class/push/push.json"
try:
old_config = json.loads(read_file(old_config_file))
except:
return
if "load_balance_push" not in old_config:
old_config["load_balance_push"] = {}
old_data = {
"push_count": task["number_rule"].get("day_num", 2),
"cycle": task["task_data"].get("cycle", "200|301|302|403|404"),
"interval": task["task_data"].get("interval", 60),
"title": task["title"],
"status": task['status'],
"module": ",".join(task["sender"])
}
for k, v in old_config["load_balance_push"].items():
if v["project"] == task["task_data"]["project"]:
v.update(old_data)
else:
old_data["project"] = task["task_data"]["project"]
old_config["load_balance_push"][int(time.time())] = old_data
write_file(old_config_file, json.dumps(old_config))
def task_config_update_hook(self, task: dict) -> None:
return self.task_config_create_hook(task)
def task_config_remove_hook(self, task: dict) -> None:
old_config_file = "/www/server/panel/class/push/push.json"
try:
old_config = json.loads(read_file(old_config_file))
except:
return
if "load_balance_push" not in old_config:
old_config["load_balance_push"] = {}
old_config["load_balance_push"] = {
k: v for k, v in old_config["load_balance_push"].items()
if v["project"] != task["task_data"]["project"]
}
def load_load_template():
if TaskTemplateConfig().get_by_id("50"):
return None
from .mods import load_task_template_by_config
load_task_template_by_config(
[{
"id": "50",
"ver": "1",
"used": True,
"source": "nginx_load_push",
"title": "负载均衡",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.load_push",
"name": "NginxLoadTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "负载名称",
"type": "select",
"default": "all",
"unit": "",
"suffix": (
"<i style='color: #999;font-style: initial;font-size: 12px;margin-right: 5px'>*</i>"
"<span style='color:#999'>选中的负载配置中,出现节点访问失败时,触发告警</span>"
),
"items": [
{
"title": "所有已配置的负载",
"value": "all"
}
]
},
{
"attr": "cycle",
"name": "成功的状态码",
"type": "textarea",
"unit": "",
"suffix": (
"<br><i style='color: #999;font-style: initial;font-size: 12px;margin-right: 5px'>*</i>"
"<span style='color:#999'>状态码以竖线分隔,如:200|301|302|403|404</span>"
),
"width": "400px",
"style": {
'height': '70px',
},
"default": "200|301|302|403|404"
}
],
"sorted": [
[
"project"
],
[
"cycle"
]
],
},
"default": {
"project": "all",
"cycle": "200|301|302|403|404"
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": False
}]
)
class ViewMsgFormat(object):
@staticmethod
def get_msg(task: dict) -> Optional[str]:
if task["template_id"] == "50":
return "<span>节点访问异常时,推送告警信息(每日推送{}次后不在推送)<span>".format(
task.get("number_rule", {}).get("day_num"))
return None
+312
View File
@@ -0,0 +1,312 @@
import json
import os
import time
from typing import Union, Optional
from .mods import TaskTemplateConfig, TaskConfig, SenderConfig, TaskRecordConfig
from .system import PushSystem
from mod.base import json_response
class PushManager:
def __init__(self):
self.template_conf = TaskTemplateConfig()
self.task_conf = TaskConfig()
self.send_config = SenderConfig()
self._send_conf_cache = {}
def _get_sender_conf(self, sender_id):
if sender_id in self._send_conf_cache:
return self._send_conf_cache[sender_id]
tmp = self.send_config.get_by_id(sender_id)
self._send_conf_cache[sender_id] = tmp
return tmp
def normalize_task_config(self, task, template) -> Union[dict, str]:
result = {}
sender = task.get("sender", None)
if sender is None:
return "未设置告警通道"
if not isinstance(sender, list):
return "告警通道设置错误"
new_sender = []
for i in sender:
sender_conf = self._get_sender_conf(i)
if not sender_conf:
continue
else:
new_sender.append(i)
if sender_conf["sender_type"] not in template["send_type_list"]:
if sender_conf["sender_type"] == "sms":
return "不支持短信告警"
return "不支持的告警方式:{}".format(sender_conf['data']["title"])
if not sender_conf["used"]:
if sender_conf["sender_type"] == "sms":
return "短信告警通道已关闭"
return "已关闭的告警方式:{}".format(sender_conf['data']["title"])
result["sender"] = new_sender
if "default" in template and template["default"]:
task_data = task.get("task_data", {})
for k, v in template["default"].items():
if k not in task_data:
task_data[k] = v
result["task_data"] = task_data
if "task_data" not in result:
result["task_data"] = {}
time_rule = task.get("time_rule", {})
if "send_interval" in time_rule:
if not isinstance(time_rule["send_interval"], int):
return "最小间隔时间设置错误"
if time_rule["send_interval"] < 0:
return "最小间隔时间设置错误"
if "time_range" in time_rule:
if not isinstance(time_rule["time_range"], list):
return "时间范围设置错误"
if not len(time_rule["time_range"]) == 2:
del time_rule["time_range"]
else:
time_range = time_rule["time_range"]
if not (isinstance(time_range[0], int) and isinstance(time_range[1], int) and
0 <= time_range[0] < time_range[1] <= 60 * 60 * 24):
return "时间范围设置错误"
result["time_rule"] = time_rule
number_rule = task.get("number_rule", {})
if "day_num" in number_rule:
if not (isinstance(number_rule["day_num"], int) and number_rule["day_num"] >= 0):
return "每日最小次数设置错误"
if "total" in number_rule:
if not (isinstance(number_rule["total"], int) and number_rule["total"] >= 0):
return "最大告警次数设置错误"
result["number_rule"] = number_rule
if "status" not in task:
result["status"] = True
if "status" in task:
if isinstance(task["status"], bool):
result["status"] = task["status"]
return result
def set_task_conf_data(self, push_data: dict) -> Optional[str]:
task_id = push_data.get("task_id", None)
template_id = push_data.get("template_id")
task = push_data.get("task_data")
target_task_conf = None
if task_id is not None:
tmp = self.task_conf.get_by_id(task_id)
if tmp is None:
target_task_conf = tmp
template = self.template_conf.get_by_id(template_id)
if not template:
return "未查询到告警模板"
if template["unique"] and not target_task_conf:
for i in self.task_conf.config:
if i["template_id"] == template["id"]:
target_task_conf = i
break
task_obj = PushSystem().get_task_object(template_id, template["load_cls"])
if not task_obj:
return "加载任务类型错误,您可以尝试修复面板"
res = self.normalize_task_config(task, template)
if isinstance(res, str):
return res
task_data = task_obj.check_task_data(res["task_data"])
if isinstance(task_data, str):
return task_data
number_rule = task_obj.check_num_rule(res["number_rule"])
if isinstance(number_rule, str):
return number_rule
time_rule = task_obj.check_time_rule(res["time_rule"])
if isinstance(time_rule, str):
return time_rule
res["task_data"] = task_data
res["number_rule"] = number_rule
res["time_rule"] = time_rule
res["keyword"] = task_obj.get_keyword(task_data)
res["source"] = task_obj.source_name
res["title"] = task_obj.get_title(task_data)
if not target_task_conf:
tmp = self.task_conf.get_by_keyword(res["source"], res["keyword"])
if tmp:
target_task_conf = tmp
if not target_task_conf:
res["id"] = self.task_conf.nwe_id()
res["template_id"] = template_id
res["status"] = True
res["pre_hook"] = {}
res["after_hook"] = {}
res["last_check"] = 0
res["last_send"] = 0
res["number_data"] = {}
res["create_time"] = time.time()
res["record_time"] = 0
self.task_conf.config.append(res)
task_obj.task_config_create_hook(res)
else:
target_task_conf.update(res)
target_task_conf["last_check"] = 0
target_task_conf["number_data"] = {} # 次数控制数据置空
task_obj.task_config_update_hook(target_task_conf)
self.task_conf.save_config()
return None
def set_task_conf(self, get):
task_id = None
try:
if hasattr(get, "task_id"):
task_id = get.task_id.strip()
if not task_id:
task_id = None
else:
self.remove_task_conf(get)
template_id = get.template_id.strip()
task = json.loads(get.task_data.strip())
except (AttributeError, json.JSONDecodeError, TypeError, ValueError):
return json_response(status=False, msg="参数错误")
push_data = {
"task_id": task_id,
"template_id": template_id,
"task_data": task,
}
res = self.set_task_conf_data(push_data)
if res:
return json_response(status=False, msg=res)
# target_task_conf = None
# if task_id is not None:
# tmp = self.task_conf.get_by_id(task_id)
# if tmp is None:
# target_task_conf = tmp
#
# template = self.template_conf.get_by_id(template_id)
# if not template:
# return json_response(status=False, msg="为查询到告警模板")
#
# if template["unique"] and not target_task_conf:
# for i in self.task_conf.config:
# if i["template_id"] == template["id"]:
# target_task_conf = i
# break
#
# task_obj = PushSystem().get_task_object(template_id, template["load_cls"])
# if not task_obj:
# return json_response(status=False, msg="加载任务类型错误,您可以尝试修复面板")
#
# res = self.normalize_task_config(task, template)
# if isinstance(res, str):
# return json_response(status=True, msg=res)
#
# task_data = task_obj.check_task_data(res["task_data"])
# if isinstance(task_data, str):
# return json_response(status=True, msg=task_data)
#
# number_rule = task_obj.check_num_rule(res["number_rule"])
# if isinstance(number_rule, str):
# return json_response(status=True, msg=number_rule)
#
# time_rule = task_obj.check_time_rule(res["time_rule"])
# if isinstance(time_rule, str):
# return json_response(status=True, msg=time_rule)
#
# res["task_data"] = task_data
# res["number_rule"] = number_rule
# res["time_rule"] = time_rule
#
# res["keyword"] = task_obj.get_keyword(task_data)
# res["source"] = task_obj.source_name
# res["title"] = task_obj.get_title(task_data)
#
# if not target_task_conf:
# tmp = self.task_conf.get_by_keyword(res["source"], res["keyword"])
# if tmp:
# target_task_conf = tmp
#
# if not target_task_conf:
# res["id"] = self.task_conf.nwe_id()
# res["template_id"] = template_id
# res["status"] = True
# res["pre_hook"] = {}
# res["after_hook"] = {}
# res["last_check"] = 0
# res["last_send"] = 0
# res["number_data"] = {}
# res["create_time"] = time.time()
# res["record_time"] = 0
# self.task_conf.config.append(res)
# task_obj.task_config_create_hook(res)
# else:
# target_task_conf.update(res)
# target_task_conf["last_check"] = 0
# target_task_conf["number_data"] = {} # 次数控制数据置空
# task_obj.task_config_update_hook(target_task_conf)
#
# self.task_conf.save_config()
return json_response(status=True, msg="告警任务保存成功")
def change_task_conf(self, get):
try:
task_id = get.task_id.strip()
except AttributeError:
return json_response(status=False, msg="参数错误")
tmp = self.task_conf.get_by_id(task_id)
if tmp is None:
return json_response(status=True, msg="为查询到告警任务")
tmp["status"] = not tmp["status"]
self.task_conf.save_config()
return json_response(status=True, msg="操作成功")
def remove_task_conf(self, get):
try:
task_id = get.task_id.strip()
except AttributeError:
return json_response(status=False, msg="参数错误")
tmp = self.task_conf.get_by_id(task_id)
if tmp is None:
return json_response(status=True, msg="为查询到告警任务")
self.task_conf.config.remove(tmp)
self.task_conf.save_config()
template = self.template_conf.get_by_id(tmp["template_id"])
if template:
task_obj = PushSystem().get_task_object(template["id"], template["load_cls"])
if task_obj:
task_obj.task_config_remove_hook(tmp)
return json_response(status=True, msg="操作成功")
@staticmethod
def clear_task_record_by_task_id(task_id):
tr_conf = TaskRecordConfig(task_id)
if os.path.exists(tr_conf.config_file_path):
os.remove(tr_conf.config_file_path)
+368
View File
@@ -0,0 +1,368 @@
# coding: utf-8
# -------------------------------------------------------------------
# 宝塔Linux面板
# -------------------------------------------------------------------
# Copyright (c) 2015-2017 宝塔软件(http:#bt.cn) All rights reserved.
# -------------------------------------------------------------------
# Author: baozi <baozi@bt.cn>
# -------------------------------------------------------------------
# 新告警的所有数据库操作
# ------------------------------
import json
import os
import types
from typing import Any, Dict, Optional, List
from threading import Lock
from uuid import uuid4
import fcntl
from .util import Sqlite, write_log, read_file, write_file
_push_db_lock = Lock()
# 代替 class/db.py 中, 离谱的两个query函数
def msg_db_query_func(self, sql, param=()):
# 执行SQL语句返回数据集
self._Sql__GetConn()
try:
return self._Sql__DB_CONN.execute(sql, self._Sql__to_tuple(param))
except Exception as ex:
return "error: " + str(ex)
def get_push_db():
db_file = "/www/server/panel/data/db/mod_push.db"
if not os.path.isdir(os.path.dirname(db_file)):
os.makedirs(os.path.dirname(db_file))
db = Sqlite()
setattr(db, "_Sql__DB_FILE", db_file)
setattr(db, "query", types.MethodType(msg_db_query_func, db))
return db
def get_table(table_name: str):
db = get_push_db()
db.table = table_name
return db
def lock_push_db():
with open("/www/server/panel/data/db/mod_push.db", mode="rb") as msg_fd:
fcntl.flock(msg_fd.fileno(), fcntl.LOCK_EX)
_push_db_lock.locked()
def unlock_push_db():
with open("/www/server/panel/data/db/mod_push.db", mode="rb") as msg_fd:
fcntl.flock(msg_fd.fileno(), fcntl.LOCK_UN)
_push_db_lock.acquire()
def push_db_locker(func):
def inner_func(*args, **kwargs):
lock_push_db()
try:
res = func(*args, **kwargs)
except Exception as e:
unlock_push_db() # 即使 报错了 也先解锁再操作
raise e
else:
unlock_push_db()
return res
return inner_func
DB_INIT_ERROR = False
PANEL_PATH = "/www/server/panel"
PUSH_DATA_PATH = "{}/data/mod_push_data".format(PANEL_PATH)
UPDATE_VERSION_FILE = "{}/update_panel.pl".format(PUSH_DATA_PATH)
UPDATE_MOD_PUSH_FILE = "{}/update_mod.pl".format(PUSH_DATA_PATH)
class BaseConfig:
config_file_path = ""
def __init__(self):
if not os.path.exists(PUSH_DATA_PATH):
os.makedirs(PUSH_DATA_PATH)
self._config: Optional[List[Dict[str, Any]]] = None
@property
def config(self) -> List[Dict[str, Any]]:
if self._config is None:
try:
self._config = json.loads(read_file(self.config_file_path))
except:
self._config = []
return self._config
def save_config(self) -> None:
write_file(self.config_file_path, json.dumps(self.config))
@staticmethod
def nwe_id() -> str:
return uuid4().hex[::2]
def get_by_id(self, target_id: str) -> Optional[Dict[str, Any]]:
for i in self.config:
if i.get("id", None) == target_id:
return i
class TaskTemplateConfig(BaseConfig):
config_file_path = "{}/task_template.json".format(PUSH_DATA_PATH)
class TaskConfig(BaseConfig):
config_file_path = "{}/task.json".format(PUSH_DATA_PATH)
def get_by_keyword(self, source: str, keyword: str) -> Optional[Dict[str, Any]]:
for i in self.config:
if i.get("source", None) == source and i.get("keyword", None) == keyword:
return i
class TaskRecordConfig(BaseConfig):
config_file_path_fmt = "%s/task_record_{}.json" % PUSH_DATA_PATH
def __init__(self, task_id: str):
super().__init__()
self.config_file_path = self.config_file_path_fmt.format(task_id)
class SenderConfig(BaseConfig):
config_file_path = "{}/sender.json".format(PUSH_DATA_PATH)
def __init__(self):
super(SenderConfig, self).__init__()
if not os.path.exists(self.config_file_path):
write_file(self.config_file_path, json.dumps([{
"id": self.nwe_id(),
"used": True,
"sender_type": "sms",
"data": {},
"original": True
}]))
def init_db():
global DB_INIT_ERROR
# id 模板id 必须唯一, 后端开发需要协商
# ver 模板版本号, 用于更新
# used 是否在使用用
# source 来源, 如Waf, rsync
# title 标题
# load_cls 要加载的类,或者从那种调用方法中获取到任务处理对象
# template 给前端,用于展示的数据
# default 默认数据,用于数据过滤, 和默认值
# unique 是否仅可唯一设置
# create_time 创建时间
create_task_template_sql = (
"CREATE TABLE IF NOT EXISTS 'task_template' ("
"'id' INTEGER PRIMARY KEY AUTOINCREMENT, "
"'ver' TEXT NOT NULL DEFAULT '1.0.0', "
"'used' INTEGER NOT NULL DEFAULT 1, "
"'source' TEXT NOT NULL DEFAULT 'site_push', "
"'title' TEXT NOT NULL DEFAULT '', "
"'load_cls' TEXT NOT NULL DEFAULT '{}', "
"'template' TEXT NOT NULL DEFAULT '{}', "
"'default' TEXT NOT NULL DEFAULT '{}', "
"'send_type_list' TEXT NOT NULL DEFAULT '[]', "
"'unique' INTEGER NOT NULL DEFAULT 0, "
"'create_time' INTEGER NOT NULL DEFAULT (strftime('%s'))"
");"
)
# source 来源, 例如waf(防火墙), rsync(文件同步)
# keyword 关键词, 不同的来源在使用中可以以此查出具体的任务,需要每个来源自己约束
# task_data 任务数据字典,字段可以自由设计
# sender 告警通道信息,为字典,可通过get_by_func字段指定从某个函数获取,用于发送
# time_rule 告警的时间规则,包含 间隔时间(send_interval), (time-range)
# number_rule 告警的次数规则,包含 每日次数(day_num), 总次数(total), 通过函数判断(get_by_func)
# status 状态是否开启
# pre_hook, after_hook 前置处理和后置处理
# record_time, 告警记录存储时间, 默认为0, 认为长时间储存
# last_check, 上次执行检查的时间
# last_send, 上次次发送时间
# number_data, 发送次数信息
create_task_sql = (
"CREATE TABLE IF NOT EXISTS 'task' ("
"'id' INTEGER PRIMARY KEY AUTOINCREMENT, "
"'template_id' INTEGER NOT NULL DEFAULT 0, "
"'source' TEXT NOT NULL DEFAULT '', "
"'keyword' TEXT NOT NULL DEFAULT '', "
"'title' TEXT NOT NULL DEFAULT '', "
"'task_data' TEXT NOT NULL DEFAULT '{}', "
"'sender' TEXT NOT NULL DEFAULT '[]', "
"'time_rule' TEXT NOT NULL DEFAULT '{}', "
"'number_rule' TEXT NOT NULL DEFAULT '{}', "
"'status' INTEGER NOT NULL DEFAULT 1, "
"'pre_hook' TEXT NOT NULL DEFAULT '{}', "
"'after_hook' TEXT NOT NULL DEFAULT '{}', "
"'last_check' INTEGER NOT NULL DEFAULT 0, "
"'last_send' INTEGER NOT NULL DEFAULT 0, "
"'number_data' TEXT NOT NULL DEFAULT '{}', "
"'create_time' INTEGER NOT NULL DEFAULT (strftime('%s')), "
"'record_time' INTEGER NOT NULL DEFAULT 0"
");"
)
create_task_record_sql = (
"CREATE TABLE IF NOT EXISTS 'task_record' ("
"'id' INTEGER PRIMARY KEY AUTOINCREMENT, "
"'template_id' INTEGER NOT NULL DEFAULT 0, "
"'task_id' INTEGER NOT NULL DEFAULT 0, "
"'do_send' TEXT NOT NULL DEFAULT '{}', "
"'send_data' TEXT NOT NULL DEFAULT '{}', "
"'result' TEXT NOT NULL DEFAULT '{}', "
"'create_time' INTEGER NOT NULL DEFAULT (strftime('%s'))"
");"
)
create_send_record_sql = (
"CREATE TABLE IF NOT EXISTS 'send_record' ("
"'id' INTEGER PRIMARY KEY AUTOINCREMENT, "
"'record_id' INTEGER NOT NULL DEFAULT 0, "
"'sender_name' TEXT NOT NULL DEFAULT '', "
"'sender_id' INTEGER NOT NULL DEFAULT 0, "
"'sender_type' TEXT NOT NULL DEFAULT '', "
"'send_data' TEXT NOT NULL DEFAULT '{}', "
"'result' TEXT NOT NULL DEFAULT '{}', "
"'create_time' INTEGER NOT NULL DEFAULT (strftime('%s'))"
");"
)
create_sender_sql = (
"CREATE TABLE IF NOT EXISTS 'sender' ("
"'id' INTEGER PRIMARY KEY AUTOINCREMENT, "
"'used' INTEGER NOT NULL DEFAULT 1, "
"'sender_type' TEXT NOT NULL DEFAULT '', "
"'name' TEXT NOT NULL DEFAULT '', "
"'data' TEXT NOT NULL DEFAULT '{}', "
"'create_time' INTEGER NOT NULL DEFAULT (strftime('%s'))"
");"
)
lock_push_db()
with get_push_db() as db:
db.execute("pragma journal_mode=wal")
res = db.execute(create_task_template_sql)
if isinstance(res, str) and res.startswith("error"):
write_log("告警系统", "task_template数据表创建错误:" + res)
DB_INIT_ERROR = True
return
res = db.execute(create_task_sql)
if isinstance(res, str) and res.startswith("error"):
write_log("告警系统", "task数据表创建错误:" + res)
DB_INIT_ERROR = True
return
res = db.execute(create_task_record_sql)
if isinstance(res, str) and res.startswith("error"):
write_log("告警系统", "task_recorde数据表创建错误:" + res)
DB_INIT_ERROR = True
return
res = db.execute(create_send_record_sql)
if isinstance(res, str) and res.startswith("error"):
write_log("告警系统", "send_record数据表创建错误:" + res)
DB_INIT_ERROR = True
return
res = db.execute(create_sender_sql)
if isinstance(res, str) and res.startswith("error"):
write_log("告警系统", "sender数据表创建错误:" + res)
DB_INIT_ERROR = True
return
db.execute(
"INSERT INTO 'sender' (id, sender_type, data) VALUES (?,?,?)",
(1, 'sms', json.dumps({"count": 0, "total": 0}))
) # 插入短信
unlock_push_db()
init_template_file = "/www/server/panel/config/mod_push_init.json"
err = load_task_template_by_file(init_template_file)
if err:
write_log("告警系统", "task_template数据表初始数据加载失败:" + res)
def _check_fields(template: dict) -> bool:
if not isinstance(template, dict):
return False
fields = ("id", "ver", "used", "source", "title", "load_cls", "template", "default", "unique", "create_time")
for field in fields:
if field not in template:
return False
return True
def load_task_template_by_config(templates: List[Dict]) -> None:
"""
通过 传入的配置信息 执行一次模板更新操作
@param templates: 模板内容,为一个数据列表
@return: 报错信息,如果返回None则表示执行成功
"""
task_template_config = TaskTemplateConfig()
add_list = []
for template in templates:
tmp = task_template_config.get_by_id(template['id'])
if tmp is not None:
tmp.update(template)
else:
add_list.append(template)
task_template_config.config.extend(add_list)
task_template_config.save_config()
# with get_table('task_template') as table:
# for template in templates:
# if not _check_fields(template):
# continue
# res = table.where("id = ?", (template['id'])).field('ver').select()
# if isinstance(res, str):
# return "数据库损坏:" + res
# if not res: # 没有就插入
# table.insert(template)
# else:
# # 版本不一致就更新版本
# if res['ver'] != template['ver']:
# template.pop("id")
# table.where("id = ?", (template['id'])).update(template)
#
def load_task_template_by_file(template_file: str) -> Optional[str]:
"""
执行一次模板更新操作
@param template_file: 模板文件路径
@return: 报错信息,如果返回None则表示执行成功
"""
if not os.path.isfile(template_file):
return "模板文件不存在,更新失败"
if DB_INIT_ERROR:
return "数据库初始化时报错,无法更新"
res = read_file(template_file)
if not isinstance(res, str):
return "数据读取失败"
try:
templates = json.loads(res)
except (json.JSONDecoder, TypeError, ValueError):
return "仅支持JSON格式数据"
if not isinstance(templates, list):
return "数据格式错误,应当为一个列表"
return load_task_template_by_config(templates)
+301
View File
@@ -0,0 +1,301 @@
import json
import os
import re
from datetime import datetime, timedelta
from typing import Tuple, Union, Optional, Iterator
from .send_tool import WxAccountMsg
from .base_task import BaseTask
from .mods import TaskTemplateConfig
from .util import read_file
def rsync_ver_is_38() -> Optional[bool]:
"""
检查rsync的版本是否为3.8。
该函数不接受任何参数。
返回值:
- None: 如果无法确定rsync的版本或文件不存在。
- bool: 如果版本确定为3.8,则返回True;否则返回False。
"""
push_file = "/www/server/panel/plugin/rsync/rsync_push.py"
if not os.path.exists(push_file):
return None
ver_info_file = "/www/server/panel/plugin/rsync/info.json"
if not os.path.exists(ver_info_file):
return None
try:
info = json.loads(read_file(ver_info_file))
except (json.JSONDecodeError, TypeError):
return None
ver = info["versions"]
ver_tuples = [int(i) for i in ver.split(".")]
if len(ver_tuples) < 3:
ver_tuples = ver_tuples.extend([0] * (3 - len(ver_tuples)))
if ver_tuples[0] < 3:
return None
if ver_tuples[1] <= 8 and ver_tuples[0] == 3:
return True
return False
class Rsync38Task(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "rsync_push"
self.template_name = "文件同步告警"
self.title = "文件同步告警"
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if "interval" not in task_data or not isinstance(task_data["interval"], int):
task_data["interval"] = 600
return task_data
def get_keyword(self, task_data: dict) -> str:
return "rsync_push"
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
has_err = self._check(task_data.get("interval", 600))
if not has_err:
return None
return {
"msg_list": [
">通知类型:文件同步告警",
">告警内容:<font color=#ff0000>文件同步执行中出错了,请及时关注文件同步情况并处理。</font> ",
]
}
@staticmethod
def _check(interval: int) -> bool:
if not isinstance(interval, int):
return False
start_time = datetime.now() - timedelta(seconds=interval * 1.2)
log_file = "{}/plugin/rsync/lsyncd.log".format("/www/server/panel")
if not os.path.exists(log_file):
return False
return LogChecker(log_file=log_file, start_time=start_time)()
def check_time_rule(self, time_rule: dict) -> Union[dict, str]:
if "send_interval" not in time_rule or not isinstance(time_rule["interval"], int):
time_rule["send_interval"] = 3 * 60
if time_rule["send_interval"] < 60:
time_rule["send_interval"] = 60
return time_rule
def filter_template(self, template: dict) -> Optional[dict]:
res = rsync_ver_is_38()
if res is None:
return None
if res:
return template
else:
return None
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "文件同步告警"
msg.msg = "同步执行出错了,请及时关注同步情况"
return msg
class Rsync39Task(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "rsync_push"
self.template_name = "文件同步告警"
self.title = "文件同步告警"
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if "interval" not in task_data or not isinstance(task_data["interval"], int):
task_data["interval"] = 600
return task_data
def get_keyword(self, task_data: dict) -> str:
return "rsync_push"
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
"""
不返回数据,以实时触发为主
"""
return None
def check_time_rule(self, time_rule: dict) -> Union[dict, str]:
if "send_interval" not in time_rule or not isinstance(time_rule["send_interval"], int):
time_rule["send_interval"] = 3 * 60
if time_rule["send_interval"] < 60:
time_rule["send_interval"] = 60
return time_rule
def filter_template(self, template: dict) -> Optional[dict]:
res = rsync_ver_is_38()
if res is None:
return None
if res is False:
return template
else:
return None
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
task_name = push_data.get("task_name", None)
msg = WxAccountMsg.new_msg()
msg.thing_type = "文件同步告警"
if task_name:
msg.msg = "文件同步任务{}出错了".format(task_name)
else:
msg.msg = "同步执行出错了,请及时关注同步情况"
return msg
class LogChecker:
"""
排序查询并获取日志内容
"""
rep_time = re.compile(r'(?P<target>(\w{3}\s+){2}(\d{1,2})\s+(\d{2}:?){3}\s+\d{4})')
format_str = '%a %b %d %H:%M:%S %Y'
err_datetime = datetime.fromtimestamp(0)
err_list = ("error", "Error", "ERROR", "exitcode = 10", "failed")
def __init__(self, log_file: str, start_time: datetime):
self.log_file = log_file
self.start_time = start_time
self.is_over_time = None # None:还没查到时间,未知, False: 可以继续网上查询, True:比较早的数据了,不再向上查询
self.has_err = False # 目前已查询的内容中是否有报错信息
def _format_time(self, log_line) -> Optional[datetime]:
try:
date_str_res = self.rep_time.search(log_line)
if date_str_res:
time_str = date_str_res.group("target")
return datetime.strptime(time_str, self.format_str)
except Exception:
return self.err_datetime
return None
# 返回日志内容
def __call__(self):
_buf = b""
file_size, fp = os.stat(self.log_file).st_size - 1, open(self.log_file, mode="rb")
fp.seek(-1, 2)
while file_size:
read_size = min(1024, file_size)
fp.seek(-read_size, 1)
buf: bytes = fp.read(read_size) + _buf
fp.seek(-read_size, 1)
if file_size > 1024:
idx = buf.find(ord("\n"))
_buf, buf = buf[:idx], buf[idx + 1:]
for i in self._get_log_line_from_buf(buf):
self._check(i)
if self.is_over_time:
return self.has_err
file_size -= read_size
return False
# 从缓冲中读取日志
@staticmethod
def _get_log_line_from_buf(buf: bytes) -> Iterator[str]:
n, m = 0, 0
buf_len = len(buf) - 1
for i in range(buf_len, -1, -1):
if buf[i] == ord("\n"):
log_line = buf[buf_len + 1 - m: buf_len - n + 1].decode("utf-8")
yield log_line
n = m = m + 1
else:
m += 1
yield buf[0: buf_len - n + 1].decode("utf-8")
# 格式化并筛选查询条件
def _check(self, log_line: str) -> None:
# 筛选日期
for err in self.err_list:
if err in log_line:
self.has_err = True
ck_time = self._format_time(log_line)
if ck_time:
self.is_over_time = self.start_time > ck_time
def load_rsync_template():
"""
加载rsync模板
"""
if TaskTemplateConfig().get_by_id("40"):
return None
from .mods import load_task_template_by_config
load_task_template_by_config(
[{
"id": "40",
"ver": "1",
"used": True,
"source": "rsync_push",
"title": "文件同步告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.rsync_push",
"name": "RsyncTask"
},
"template": {
"field": [
],
"sorted": [
]
},
"default": {
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": True
}]
)
RsyncTask = Rsync39Task
if rsync_ver_is_38() is True:
RsyncTask = Rsync38Task
def push_rsync_by_task_name(task_name: str):
from .system import push_by_task_keyword
push_data = {
"task_name": task_name,
"msg_list": [
">通知类型:文件同步告警",
">告警内容:<font color=#ff0000>文件同步任务{}在执行中出错了,请及时关注文件同步情况并处理。</font> ".format(
task_name),
]
}
push_by_task_keyword("rsync_push", "rsync_push", push_data=push_data)
class ViewMsgFormat(object):
@staticmethod
def get_msg(task: dict) -> Optional[str]:
if task["template_id"] == "40":
return "<span>文件同步出现异常时,推送告警信息(每日推送{}次后不在推送)<span>".format(
task.get("number_rule", {}).get("day_num"))
return None
+133
View File
@@ -0,0 +1,133 @@
import ipaddress
import re
from .util import get_config_value
class WxAccountMsgBase:
@classmethod
def new_msg(cls):
return cls()
def set_ip_address(self, server_ip, local_ip):
pass
def to_send_data(self):
return "", {}
class WxAccountMsg(WxAccountMsgBase):
def __init__(self):
self.ip_address: str = ""
self.thing_type: str = ""
self.msg: str = ""
self.next_msg: str = ""
def set_ip_address(self, server_ip, local_ip):
self.ip_address = "{}({})".format(server_ip, local_ip)
if len(self.ip_address) > 32:
self.ip_address = self.ip_address[:29] + "..."
def to_send_data(self):
res = {
"first": {},
"keyword1": {
"value": self.ip_address,
},
"keyword2": {
"value": self.thing_type,
},
"keyword3": {
"value": self.msg,
}
}
if self.next_msg != "":
res["keyword4"] = {"value": self.next_msg}
return "", res
class WxAccountLoginMsg(WxAccountMsgBase):
tid = "RJNG8dBZ5Tb9EK6j6gOlcAgGs2Fjn5Fb07vZIsYg1P4"
def __init__(self):
self.login_name: str = ""
self.login_ip: str = ""
self.thing_type: str = ""
self.login_type: str = ""
self.address: str = ""
self._server_name: str = ""
def set_ip_address(self, server_ip, local_ip):
if self._server_name == "":
self._server_name = "服务器IP{}".format(server_ip)
def _get_server_name(self):
data = get_config_value("title") # 若获得别名,则使用别名
if data != "":
self._server_name = data
def to_send_data(self):
self._get_server_name()
if self.address.startswith(">归属地:"):
self.address = self.address[5:]
if self.address == "":
self.address = "未知的归属地"
if not _is_ipv4(self.login_ip):
self.login_ip = "ipv6-can not show"
res = {
"thing10": {
"value": self._server_name,
},
"character_string9": {
"value": self.login_ip,
},
"thing7": {
"value": self.login_type,
},
"thing11": {
"value": self.address,
},
"thing2": {
"value": self.login_name,
}
}
return self.tid, res
# 处理短信告警信息的不规范问题
def sms_msg_normalize(sm_args: dict) -> dict:
for key, val in sm_args.items():
sm_args[key] = _norm_sms_push_argv(str(val))
return sm_args
def _norm_sms_push_argv(data):
"""
@处理短信参数,否则会被拦截
"""
if _is_ipv4(data):
tmp1 = data.split('.')
return '{}_***_***_{}'.format(tmp1[0], tmp1[3])
data = data.replace(".", "_").replace("+", "+")
return data
def _is_ipv4(data: str) -> bool:
try:
ipaddress.IPv4Address(data)
except:
return False
return True
def _is_domain(domain):
rep_domain = re.compile(r"^([\w\-*]{1,100}\.){1,10}([\w\-]{1,24}|[\w\-]{1,24}\.[\w\-]{1,24})$")
if rep_domain.match(domain):
return True
return False
File diff suppressed because it is too large Load Diff
+571
View File
@@ -0,0 +1,571 @@
[
{
"id": "1",
"ver": "1",
"used": true,
"source": "site_ssl",
"title": "网站证书(SSL)到期",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "SSLTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "网站",
"type": "select",
"default": "all",
"items": [
{
"title": "所有网站",
"value": "all"
}
]
},
{
"attr": "cycle",
"name": "剩余天数",
"type": "number",
"suffix": "",
"unit": "天",
"default": 15
}
],
"sorted": [
[
"project"
],
[
"cycle"
]
]
},
"default": {
"project": "all",
"cycle": 15
},
"advanced_default": {
"number_rule": {
"total": 2
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": false
},
{
"id": "2",
"ver": "1",
"used": true,
"source": "site_end_time",
"title": "网站到期",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "SiteEndTimeTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "剩余天数",
"type": "number",
"unit": "天",
"suffix": "",
"default": 7
}
],
"sorted": [
[
"cycle"
]
]
},
"default": {
"cycle": 7
},
"advanced_default": {
"number_rule": {
"total": 2
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
},
{
"id": "3",
"ver": "1",
"used": true,
"source": "panel_pwd_end_time",
"title": "面板密码有效期",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "PanelPwdEndTimeTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "剩余天数",
"type": "number",
"unit": "天",
"suffix": "",
"default": 15
}
],
"sorted": [
[
"cycle"
]
]
},
"default": {
"cycle": 15
},
"advanced_default": {
"number_rule": {
"total": 2
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
},
{
"id": "4",
"ver": "1",
"used": true,
"source": "ssh_login_error",
"title": "SSH登录失败告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "SSHLoginErrorTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "触发条件",
"type": "number",
"unit": "分钟",
"suffix": "内,",
"default": 30
},
{
"attr": "count",
"name": "登录失败",
"type": "number",
"unit": "次",
"suffix": "",
"default": 3
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"cycle",
"count"
],
[
"interval"
]
]
},
"default": {
"cycle": 30,
"count": 3,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
},
"time_rule": {
"send_interval": 600
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
},
{
"id": "5",
"ver": "1",
"used": true,
"source": "services",
"title": "服务停止告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "ServicesTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "通知类型",
"type": "select",
"default": null,
"items": [
]
},
{
"attr": "count",
"name": "自动重启",
"type": "radio",
"suffix": "",
"default": 1,
"items": [
{
"title": "自动尝试重启项目",
"value": 1
},
{
"title": "不做重启尝试",
"value": 2
}
]
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"project"
],
[
"count"
],
[
"interval"
]
]
},
"default": {
"project": "",
"count": 2,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": false
},
{
"id": "6",
"ver": "1",
"used": true,
"source": "panel_safe_push",
"title": "面板安全告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "PanelSafePushTask"
},
"template": {
"field": [
{
"attr": "help",
"name": "告警内容",
"type": "help",
"unit": "",
"style": {
"margin-top": "6px"
},
"list": [
"面板用户变更、面板日志删除、面板开启开发者"
],
"suffix": "",
"default": 600
}
],
"sorted": [
[
"help"
]
]
},
"default": {
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
},
{
"id": "7",
"ver": "1",
"used": true,
"source": "ssh_login",
"title": "SSH登录告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "SSHLoginTask"
},
"template": {
"field": [
],
"sorted": [
]
},
"default": {
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
},
{
"id": "8",
"ver": "1",
"used": true,
"source": "panel_login",
"title": "面板登录告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "PanelLoginTask"
},
"template": {
"field": [
],
"sorted": [
]
},
"default": {
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": true
},
{
"id": "9",
"ver": "1",
"used": true,
"source": "project_status",
"title": "项目停止告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "ProjectStatusTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "项目类型",
"type": "select",
"default": 1,
"items": [
{
"title": "Java项目",
"value": 1
},
{
"title": "Node项目",
"value": 2
},
{
"title": "Go项目",
"value": 3
},
{
"title": "Python项目",
"value": 4
},
{
"title": "其他项目",
"value": 5
}
]
},
{
"attr": "project",
"name": "项目名称",
"type": "select",
"default": null,
"all_items": null,
"items": [
]
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
},
{
"attr": "count",
"name": "自动重启",
"type": "radio",
"suffix": "",
"default": 1,
"items": [
{
"title": "自动尝试重启项目",
"value": 1
},
{
"title": "不做重启尝试",
"value": 2
}
]
}
],
"sorted": [
[
"cycle"
],
[
"project"
],
[
"interval"
],
[
"count"
]
]
},
"default": {
"cycle": 1,
"project": "",
"interval": 600,
"count": 2
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": false
},
{
"id": "10",
"ver": "1",
"used": true,
"source": "panel_update",
"title": "面板更新提醒",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.site_push",
"name": "PanelUpdateTask"
},
"template": {
"field": [
],
"sorted": [
]
},
"default": {
},
"advanced_default": {
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": true
}
]
+422
View File
@@ -0,0 +1,422 @@
import os
import time
from typing import Optional, List, Tuple, Dict, Type, Any, Union
import datetime
from threading import Thread
from .base_task import BaseTask
from .mods import TaskTemplateConfig, TaskConfig, TaskRecordConfig, SenderConfig
from .send_tool import sms_msg_normalize
from .tool import load_task_cls_by_path, load_task_cls_by_function, T_CLS
from .util import get_server_ip, get_network_ip, format_date, get_config_value
from .compatible import rsync_compatible
WAIT_TASK_LIST: List[Thread] = []
class PushSystem:
def __init__(self):
self.task_cls_cache: Dict[str, Type[T_CLS]] = {}
self._today_zero: Optional[datetime.datetime] = None
self._sender_type_class: Optional[dict] = {}
self.sd_cfg = SenderConfig()
def sender_cls(self, sender_type: str):
if not self._sender_type_class:
from mod.base.msg import WeiXinMsg, MailMsg, WebHookMsg, FeiShuMsg, DingDingMsg, SMSMsg, WeChatAccountMsg
self._sender_type_class = {
"weixin": WeiXinMsg,
"mail": MailMsg,
"webhook": WebHookMsg,
"feishu": FeiShuMsg,
"dingding": DingDingMsg,
"sms": SMSMsg,
"wx_account": WeChatAccountMsg,
}
return self._sender_type_class[sender_type]
@staticmethod
def can_run_task_list() -> Tuple[List[dict], Dict[int, dict]]:
result = []
result_template = {}
task_template_ids = set()
for task in TaskConfig().config:
if not task["status"]:
continue
task_template_ids.add(task['template_id'])
# 间隔检测时间未到跳过
if "interval" in task["task_data"] and isinstance(task["task_data"]["interval"], int):
if time.time() < task["last_check"] + task["task_data"]["interval"]:
continue
result.append(task)
for template in TaskTemplateConfig().config:
if template["id"] not in task_template_ids:
continue
result_template[template['id']] = template
return result, result_template
def get_task_object(self, template_id, load_cls_data: dict) -> Optional[BaseTask]:
if template_id in self.task_cls_cache:
return self.task_cls_cache[template_id]()
if "load_type" not in load_cls_data:
return None
if load_cls_data["load_type"] == "func":
cls = load_task_cls_by_function(
name=load_cls_data["name"],
func_name=load_cls_data["func_name"],
is_model=load_cls_data.get("is_model", False),
model_index=load_cls_data.get("is_model", ''),
args=load_cls_data.get("args", None),
sub_name=load_cls_data.get("sub_name", None),
)
else:
cls_path = load_cls_data["cls_path"]
cls = load_task_cls_by_path(cls_path, load_cls_data["name"])
if not cls:
return None
self.task_cls_cache[template_id] = cls
return cls()
def run(self):
rsync_compatible()
task_list, task_template = self.can_run_task_list()
for t in task_list:
if t["template_id"] not in task_template:
continue
template = task_template[t["template_id"]]
if not template["used"]:
continue
print(t)
print(_PushRunner(t, template, self)())
print(">>>>>>>>>>>>>>>>>>>>>>>>>>>>")
global WAIT_TASK_LIST
if WAIT_TASK_LIST: # 有任务启用子线程的,要等到这个线程结束,再结束主线程
for i in WAIT_TASK_LIST:
i.join()
def get_today_zero(self) -> datetime.datetime:
if self._today_zero is None:
t = datetime.datetime.today()
t_zero = datetime.datetime.combine(t, datetime.time.min)
self._today_zero = t_zero
return self._today_zero
class _PushRunner:
def __init__(self, task: dict, template: dict, push_system: PushSystem, custom_push_data: Optional[dict] = None):
self._public_push_data: Optional[dict] = None
self.result: dict = {
"do_send": False,
"stop_msg": "",
"push_data": {},
"check_res": False,
"check_stop_on": "",
"send_data": {},
} # 记录结果
self.change_fields = set() # 记录task变化值
self.task_obj: Optional[BaseTask] = None
self.task = task
self.template = template
self.push_system = push_system
self._add_hook_msg: Optional[str] = None # 记录前置钩子处理后的追加信息
self.custom_push_data = custom_push_data
self.tr_cfg = TaskRecordConfig(task["id"])
self.is_number_rule_by_func = False # 记录这个任务是否使用自定义的次数检测, 如果是,就不需要做次数更新
def save_result(self):
t = TaskConfig()
tmp = t.get_by_id(self.task["id"])
if tmp:
for f in self.change_fields:
tmp[f] = self.task[f]
if self.result["do_send"]:
tmp["last_send"] = int(time.time())
tmp["last_check"] = int(time.time())
t.save_config()
if self.result["push_data"]:
result_data = self.result.copy()
self.tr_cfg.config.append(
{
"id": self.tr_cfg.nwe_id(),
"template_id": self.template["id"],
"task_id": self.task["id"],
"do_send": result_data.pop("do_send"),
"send_data": result_data.pop("push_data"),
"result": result_data,
"create_time": int(time.time()),
}
)
self.tr_cfg.save_config()
@property
def public_push_data(self) -> dict:
if self._public_push_data is None:
self._public_push_data = {
'ip': get_server_ip(),
'local_ip': get_network_ip(),
'server_name': get_config_value('title')
}
data = self._public_push_data.copy()
data['time'] = format_date()
data['timestamp'] = int(time.time())
return data
def __call__(self):
self.run()
self.save_result()
if self.task_obj:
self.task_obj.task_run_end_hook(self.result)
return self.result_to_return()
def result_to_return(self) -> dict:
return self.result
def run(self):
self.task_obj = self.push_system.get_task_object(self.template["id"], self.template["load_cls"])
if not self.task_obj:
self.result["stop_msg"] = "任务类加载失败"
return
if self.custom_push_data is None:
push_data = self.task_obj.get_push_data(self.task["id"], self.task["task_data"])
if not push_data:
return
else:
push_data = self.custom_push_data
self.result["push_data"] = push_data
# 执行前置钩子
if self.task["pre_hook"] and "hook_type" in self.task["pre_hook"]:
if not self.run_hook(self.task["pre_hook"], "pre_hook"):
return
# 执行时间规则判断
if not self.run_time_rule(self.task["time_rule"]):
return
# 执行时间规则判断
if not self.number_rule(self.task["number_rule"]):
return
# 执行发送信息
self.send_message(push_data)
self.change_fields.add("number_data")
if "day_num" not in self.task["number_data"]:
self.task["number_data"]["day_num"] = 0
if "total" not in self.task["number_data"]:
self.task["number_data"]["total"] = 0
self.task["number_data"]["day_num"] += 1
self.task["number_data"]["total"] += 1
self.task["number_data"]["time"] = int(time.time())
# 执行后置钩子
if self.task["after_hook"] and "hook_type" in self.task["after_hook"]:
self.run_hook(self.task["after_hook"], "after_hook")
# todo: 下个版本实现一些自定义的hook函数,同时实现用户脚本的hook记录在 self.result 最后统一储存
def run_hook(self, hook_data: dict, hook_name: str) -> bool:
"""
执行hook操作,并返回是否继续执行, 并将hook的执行结果记录
@param hook_name: 钩子的名称,如:after_hook, pre_hook
@param hook_data: 执行的内容
@return:
"""
return True
def run_time_rule(self, time_rule: dict) -> bool:
if "send_interval" in time_rule and time_rule["send_interval"] > 0:
if self.task["last_send"] + time_rule["send_interval"] > time.time():
self.result['stop_msg'] = '小于最小发送时间,不进行发送'
self.result['check_stop_on'] = "time_rule_send_interval"
return False
time_range = time_rule.get("time_range", None)
if time_range and isinstance(time_range, list) and len(time_range) == 2:
t_zero = self.push_system.get_today_zero()
start_time = t_zero + datetime.timedelta(seconds=time_range[0])
end_time = t_zero + datetime.timedelta(seconds=time_range[1])
if not start_time < datetime.datetime.now() < end_time:
self.result['stop_msg'] = '不在可发送告警的时间范围之内'
self.result['check_stop_on'] = "time_rule_time_range"
return False
return True
def number_rule(self, number_rule: dict) -> bool:
number_data = self.task.get("number_data", {})
# 判断通过 自定义函数的方式确认是否达到发送次数
if "get_by_func" in number_rule and isinstance(number_rule["get_by_func"], str):
f = getattr(self.task_obj, number_rule["get_by_func"], None)
if f is not None and callable(f):
res = f(self.task["id"], self.task["task_data"], number_data, self.result["push_data"])
if isinstance(res, str):
self.result['stop_msg'] = res
self.result['check_stop_on'] = "number_rule_get_by_func"
return False
# 只要是走了使用函数检查的,不再处理默认情况 change_fields 中不添加 number_data
return True
if "day_num" in number_rule and isinstance(number_rule["day_num"], int) and number_rule["day_num"] > 0:
record_time = number_data.get("time", 0)
if record_time < self.push_system.get_today_zero().timestamp(): # 昨日触发
self.task["number_data"]["day_num"] = record_num = 0
self.task["number_data"]["time"] = time.time()
self.change_fields.add("number_data")
else:
record_num = self.task["number_data"].get("day_num")
if record_num >= number_rule["day_num"]:
self.result['stop_msg'] = "超过每日限制次数:{}".format(number_rule["day_num"])
self.result['check_stop_on'] = "number_rule_day_num"
return False
if "total" in number_rule and isinstance(number_rule["total"], int) and number_rule["total"] > 0:
record_total = number_data.get("total", 0)
if record_total >= number_rule["total"]:
self.result['stop_msg'] = "超过最大限制次数:{}".format(number_rule["total"])
self.result['check_stop_on'] = "number_rule_total"
return False
return True
def send_message(self, push_data: dict):
self.result["do_send"] = True
self.result["push_data"] = push_data
wx_account = []
for sender_id in self.task["sender"]:
conf = self.push_system.sd_cfg.get_by_id(sender_id)
if conf is None:
continue
if not conf["used"]:
self.result["send_data"][sender_id] = "告警通道{}已关闭,跳过发送".format(conf["data"].get("title"))
continue
sd_cls = self.push_system.sender_cls(conf["sender_type"])
if conf["sender_type"] == "weixin":
res = sd_cls(conf).send_msg(
self.task_obj.to_weixin_msg(push_data, self.public_push_data),
self.task_obj.title
)
elif conf["sender_type"] == "mail":
res = sd_cls(conf).send_msg(
self.task_obj.to_mail_msg(push_data, self.public_push_data),
self.task_obj.title
)
elif conf["sender_type"] == "webhook":
res = sd_cls(conf).send_msg(
self.task_obj.to_web_hook_msg(push_data, self.public_push_data),
self.task_obj.title,
self.task_obj.title
)
elif conf["sender_type"] == "feishu":
res = sd_cls(conf).send_msg(
self.task_obj.to_feishu_msg(push_data, self.public_push_data),
self.task_obj.title
)
elif conf["sender_type"] == "dingding":
res = sd_cls(conf).send_msg(
self.task_obj.to_dingding_msg(push_data, self.public_push_data),
self.task_obj.title
)
elif conf["sender_type"] == "sms":
sm_type, sm_args = self.task_obj.to_sms_msg(push_data, self.public_push_data)
if not sm_type or not sm_args:
continue
sm_args = sms_msg_normalize(sm_args)
res = sd_cls(conf).send_msg(sm_type, sm_args)
elif conf["sender_type"] == "wx_account":
wx_account.append(conf)
continue
else:
continue
if isinstance(res, str) and res.find("Traceback") != -1:
self.result["send_data"][sender_id] = "执行信息发送过程中报错了, 未发送成功"
if isinstance(res, str):
self.result["send_data"][sender_id] = res
else:
self.result["send_data"][sender_id] = 1
if len(wx_account) > 0:
sd_cls = self.push_system.sender_cls("wx_account")
res = sd_cls(*wx_account).send_msg(self.task_obj.to_wx_account_msg(push_data, self.public_push_data))
for i in wx_account:
if isinstance(res, str):
self.result["send_data"][i["id"]] = res
else:
self.result["send_data"][i["id"]] = 1
def push_by_task_keyword(source: str, keyword: str, push_data: Optional[dict] = None) -> Union[str, dict]:
"""
通过关键字查询告警任务,并发送信息
@param push_data:
@param source:
@type keyword:
@return:
"""
push_system = PushSystem()
target_task = None
for i in TaskConfig().config:
if i["source"] == source and i["keyword"] == keyword:
target_task = i
break
if not target_task:
return "未查找到该任务"
target_template = TaskTemplateConfig().get_by_id(target_task["template_id"])
if not target_template["used"]:
return "该任务类型已被禁止使用"
if not target_task['status']:
return "该任务已停止"
return _PushRunner(target_task, target_template, push_system, push_data)()
def push_by_task_id(task_id: str, push_data: Optional[dict] = None):
"""
通过任务id触发告警 并 发送信息
@param push_data:
@param task_id:
@return:
"""
push_system = PushSystem()
target_task = TaskConfig().get_by_id(task_id)
if not target_task:
return "未查找到该任务"
target_template = TaskTemplateConfig().get_by_id(target_task["template_id"])
if not target_template["used"]:
return "该任务类型已被禁止使用"
if not target_task['status']:
return "该任务已停止"
return _PushRunner(target_task, target_template, push_system, push_data)()
def get_push_public_data():
data = {
'ip': get_server_ip(),
'local_ip': get_network_ip(),
'server_name': get_config_value('title'),
'time': format_date(),
'timestamp': int(time.time())}
return data
+421
View File
@@ -0,0 +1,421 @@
import json
import os
import sys
import threading
import time
from datetime import datetime, timedelta
from importlib import import_module
from typing import Tuple, Union, Optional, List
import psutil
from .send_tool import WxAccountMsg
from .base_task import BaseTask
from .mods import PUSH_DATA_PATH
from .util import read_file, write_file, get_config_value
from .system import WAIT_TASK_LIST
try:
if "/www/server/panel/class" not in sys.path:
sys.path.insert(0, "/www/server/panel/class")
from panel_msg.collector import SitePushMsgCollect, SystemPushMsgCollect
except ImportError:
SitePushMsgCollect = None
SystemPushMsgCollect = None
def _get_panel_name() -> str:
data = get_config_value("title") # 若获得别名,则使用别名
if data == "":
data = "宝塔面板"
return data
class PanelSysDiskTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "system_disk"
self.template_name = "首页磁盘告警"
self.wx_msg = ""
def get_title(self, task_data: dict) -> str:
return "挂载目录【{}】的磁盘余量告警".format(task_data["project"])
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if task_data["project"] not in [i[0] for i in self._get_disk_name()]:
return "指定的磁盘不存在"
if not (isinstance(task_data['cycle'], int) and task_data['cycle'] in (1, 2)):
return "类型参数错误"
if not (isinstance(task_data['count'], int) and task_data['count'] >= 1):
return "阈值参数错误"
if task_data['cycle'] == 2 and task_data['count'] >= 100:
return "阈值参数错误, 设置的检查范围不正确"
task_data['interval'] = 600
return task_data
@staticmethod
def _get_disk_name() -> list:
"""获取硬盘挂载点"""
if "/www/server/panel" not in sys.path:
sys.path.insert(0, "/www/server/panel")
system_modul = import_module('.system', package="class")
system = getattr(system_modul, "system")
disk_info = system.GetDiskInfo2(None, human=False)
return [(d.get("path"), d.get("size")[0]) for d in disk_info]
@staticmethod
def _get_disk_info() -> list:
"""获取硬盘挂载点"""
if "/www/server/panel" not in sys.path:
sys.path.insert(0, "/www/server/panel")
system_modul = import_module('.system', package="class")
system = getattr(system_modul, "system")
disk_info = system.GetDiskInfo2(None, human=False)
return disk_info
def get_keyword(self, task_data: dict) -> str:
print(task_data)
return task_data["project"]
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
disk_info = self._get_disk_info()
unsafe_disk_list = []
for d in disk_info:
if task_data["project"] != d["path"]:
continue
free = int(d["size"][2]) / 1048576
proportion = int(d["size"][3] if d["size"][3][-1] != "%" else d["size"][3][:-1])
if task_data["cycle"] == 1 and free < task_data["count"]:
unsafe_disk_list.append(
"挂载在【{}】上的磁盘剩余容量为{}G,小于告警值{}G.".format(
d["path"], round(free, 2), task_data["count"])
)
self.wx_msg = "剩余容量小于{}G".format(task_data["count"])
elif task_data["cycle"] == 2 and proportion > task_data["count"]:
unsafe_disk_list.append(
"挂载在【{}】上的磁盘已使用容量为{}%,大于告警值{}%.".format(
d["path"], round(proportion, 2), task_data["count"])
)
self.wx_msg = "占用量大于{}%".format(task_data["count"])
if len(unsafe_disk_list) == 0:
return None
return {
"msg_list": [
">通知类型:磁盘余量告警",
">告警内容:\n" + "\n".join(unsafe_disk_list)
]
}
def filter_template(self, template: dict) -> Optional[dict]:
for (path, total_size) in self._get_disk_name():
template["field"][0]["items"].append({
"title": "【{}】的磁盘".format(path),
"value": path,
"count_default": round((int(total_size) * 0.2) / 1024 / 1024, 1)
})
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return 'machine_exception|磁盘余量告警', {
'name': _get_panel_name(),
'type': "磁盘空间不足",
}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "宝塔首页磁盘告警"
if len(self.wx_msg) > 20:
self.wx_msg = self.wx_msg[:17] + "..."
msg.msg = self.wx_msg
return msg
class PanelSysCPUTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "system_cpu"
self.template_name = "首页CPU告警"
self.title = "首页CPU告警"
self.cpu_count = 0
self._tip_file = "{}/system_cpu.tip".format(PUSH_DATA_PATH)
self._tip_data: Optional[List[Tuple[float, float]]] = None
@property
def cache_list(self) -> List[Tuple[float, float]]:
if self._tip_data is not None:
return self._tip_data
try:
self._tip_data = json.loads(read_file(self._tip_file))
except:
self._tip_data = []
return self._tip_data
def save_cache_list(self):
write_file(self._tip_file, json.dumps(self.cache_list))
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not (isinstance(task_data['cycle'], int) and task_data['cycle'] >= 1):
return "时间参数错误"
if not (isinstance(task_data['count'], int) and task_data['count'] >= 1):
return "阈值参数错误,至少为1%"
task_data['interval'] = 60
return task_data
def get_keyword(self, task_data: dict) -> str:
return "system_cpu"
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
expiration = datetime.now() - timedelta(seconds=task_data["cycle"] * 60 + 10)
for i in range(len(self.cache_list) - 1, -1, -1):
data_time, _ = self.cache_list[i]
if datetime.fromtimestamp(data_time) < expiration:
del self.cache_list[i]
# 记录下次的
def thread_get_cpu_data():
self.cache_list.append((time.time(), psutil.cpu_percent(10)))
self.save_cache_list()
thread_active = threading.Thread(target=thread_get_cpu_data, args=())
thread_active.start()
WAIT_TASK_LIST.append(thread_active)
if len(self.cache_list) < task_data["cycle"]: # 小于指定次数不推送
return None
if len(self.cache_list) > 0:
avg_data = sum(i[1] for i in self.cache_list) / len(self.cache_list)
else:
avg_data = 0
if avg_data < task_data["count"]:
return None
else:
self.cache_list.clear()
self.cpu_count = round(avg_data, 2)
s_list = [
">通知类型:CPU高占用告警",
">告警内容:最近{}分钟内机器CPU平均占用率为{}%,高于告警值{}%".format(
task_data["cycle"], round(avg_data, 2), task_data["count"]),
]
return {
"msg_list": s_list,
}
def filter_template(self, template: dict) -> Optional[dict]:
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return 'machine_exception|CPU高占用告警', {
'name': _get_panel_name(),
'type': "CPU高占用",
}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "宝塔首页cpu告警"
msg.msg = "主机CPU占用超过:{}%".format(self.cpu_count)
msg.next_msg = "请登录面板,查看主机情况"
return msg
class PanelSysLoadTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "system_load"
self.template_name = "首页负载告警"
self.title = "首页负载告警"
self.avg_data = 0
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not (isinstance(task_data['cycle'], int) and task_data['cycle'] >= 1):
return "时间参数错误"
if not (isinstance(task_data['count'], int) and task_data['count'] >= 1):
return "阈值参数错误,至少为1%"
task_data['interval'] = 60 * task_data['cycle']
return task_data
def get_keyword(self, task_data: dict) -> str:
return "system_load"
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
now_load = os.getloadavg()
cpu_count = psutil.cpu_count()
now_load = [i / (cpu_count * 2) * 100 for i in now_load]
need_push = False
avg_data = 0
if task_data["cycle"] == 15 and task_data["count"] < now_load[2]:
avg_data = now_load[2]
need_push = True
elif task_data["cycle"] == 5 and task_data["count"] < now_load[1]:
avg_data = now_load[1]
need_push = True
elif task_data["cycle"] == 1 and task_data["count"] < now_load[0]:
avg_data = now_load[0]
need_push = True
if need_push is False:
return None
self.avg_data = avg_data
return {
"msg_list": [
">通知类型:负载超标告警",
">告警内容:最近{}分钟内机器平均负载率为{}%,高于{}%告警值".format(
task_data["cycle"], round(avg_data, 2), task_data["count"]),
]
}
def filter_template(self, template: dict) -> Optional[dict]:
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return 'machine_exception|负载超标告警', {
'name': _get_panel_name(),
'type': "平均负载过高",
}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "宝塔首页负载告警"
msg.msg = "主机负载超过:{}%".format(round(self.avg_data, 2))
msg.next_msg = "请登录面板,查看主机情况"
return msg
class PanelSysMEMTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "system_mem"
self.template_name = "首页内存告警"
self.title = "首页内存告警"
self.wx_data = 0
self._tip_file = "{}/system_mem.tip".format(PUSH_DATA_PATH)
self._tip_data: Optional[List[Tuple[float, float]]] = None
@property
def cache_list(self) -> List[Tuple[float, float]]:
if self._tip_data is not None:
return self._tip_data
try:
self._tip_data = json.loads(read_file(self._tip_file))
except:
self._tip_data = []
return self._tip_data
def save_cache_list(self):
write_file(self._tip_file, json.dumps(self.cache_list))
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not (isinstance(task_data['cycle'], int) and task_data['cycle'] >= 1):
return "次数参数错误"
if not (isinstance(task_data['count'], int) and task_data['count'] >= 1):
return "阈值参数错误,至少为1%"
task_data['interval'] = task_data['cycle'] * 60
return task_data
def get_keyword(self, task_data: dict) -> str:
return "system_mem"
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
mem = psutil.virtual_memory()
real_used: float = (mem.total - mem.free - mem.buffers - mem.cached) / mem.total
stime = datetime.now()
expiration = stime - timedelta(seconds=task_data["cycle"] * 60 + 10)
self.cache_list.append((stime.timestamp(), real_used))
for i in range(len(self.cache_list) - 1, -1, -1):
data_time, _ = self.cache_list[i]
if datetime.fromtimestamp(data_time) < expiration:
del self.cache_list[i]
avg_data = sum(i[1] for i in self.cache_list) / len(self.cache_list)
if avg_data * 100 < task_data["count"]:
self.save_cache_list()
return None
else:
self.cache_list.clear()
self.save_cache_list()
self.wx_data = round(avg_data * 100, 2)
return {
'msg_list': [
">通知类型:内存高占用告警",
">告警内容:最近{}分钟内机器内存平均占用率为{}%,高于告警值{}%".format(
task_data["cycle"], round(avg_data * 100, 2), task_data["count"]),
]
}
def filter_template(self, template: dict) -> Optional[dict]:
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return 'machine_exception|内存高占用告警', {
'name': _get_panel_name(),
'type': "内存高占用",
}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "宝塔首页内存告警"
msg.msg = "主机内存占用超过:{}%".format(self.wx_data)
msg.next_msg = "请登录面板,查看主机情况"
return msg
class ViewMsgFormat(object):
_FORMAT = {
"20": (
lambda x: "<span>挂载在{}上的磁盘{}触发</span>".format(
x.get("project"),
"余量不足%.1fG" % round(x.get("count"), 1) if x.get("cycle") == 1 else "占用超过%d%%" % x.get("count"),
)
),
"21": (
lambda x: "<span>{}分钟内平均CUP占用超过{}%触发</span>".format(
x.get("cycle"), x.get("count")
)
),
"22": (
lambda x: "<span>{}分钟内平均负载超过{}%触发</span>".format(
x.get("cycle"), x.get("count")
)
),
"23": (
lambda x: "<span>{}分钟内内存使用率超过{}%触发</span>".format(
x.get("cycle"), x.get("count")
)
)
}
def get_msg(self, task: dict) -> Optional[str]:
if task["template_id"] in self._FORMAT:
return self._FORMAT[task["template_id"]](task["task_data"])
return None
+284
View File
@@ -0,0 +1,284 @@
[
{
"id": "20",
"ver": "1",
"used": true,
"source": "system_disk",
"title": "首页磁盘告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.system_push",
"name": "PanelSysDiskTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "磁盘信息",
"type": "select",
"items": [
]
},
{
"attr": "cycle",
"name": "检测类型",
"type": "radio",
"suffix": "",
"default": 2,
"items": [
{
"title": "剩余容量",
"value": 1
},
{
"title": "占用百分比",
"value": 2
}
]
},
{
"attr": "count",
"name": "占用率超过",
"type": "number",
"unit": "%",
"suffix": "后触发告警",
"default": 80,
"err_msg_prefix": "磁盘阈值"
}
],
"sorted": [
[
"project"
],
[
"cycle"
],
[
"count"
]
]
},
"default": {
"project": "/",
"cycle": 2,
"count": 80
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": false
},
{
"id": "21",
"ver": "1",
"used": true,
"source": "system_cpu",
"title": "首页CPU告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.system_push",
"name": "PanelSysCPUTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "每",
"type": "select",
"unit": "分钟",
"suffix": "内平均",
"width": "70px",
"disabled": true,
"default": 5,
"items": [
{
"title": "1",
"value": 3
},
{
"title": "5",
"value": 5
},
{
"title": "15",
"value": 15
}
]
},
{
"attr": "count",
"name": "CPU占用超过",
"type": "number",
"unit": "%",
"suffix": "后触发告警",
"default": 80,
"err_msg_prefix": "CPU"
}
],
"sorted": [
[
"cycle",
"count"
]
]
},
"default": {
"cycle": 5,
"count": 80
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": true
},
{
"id": "22",
"ver": "1",
"used": true,
"source": "system_load",
"title": "首页负载告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.system_push",
"name": "PanelSysLoadTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "每",
"type": "select",
"unit": "分钟",
"suffix": "内平均",
"default": 5,
"width": "70px",
"disabled": true,
"items": [
{
"title": "1",
"value": 1
},
{
"title": "5",
"value": 5
},
{
"title": "15",
"value": 15
}
]
},
{
"attr": "count",
"name": "负载超过",
"type": "number",
"unit": "%",
"suffix": "后触发告警",
"default": 80,
"err_msg_prefix": "负载"
}
],
"sorted": [
[
"cycle",
"count"
]
]
},
"default": {
"cycle": 5,
"count": 80
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": true
},
{
"id": "23",
"ver": "1",
"used": true,
"source": "system_mem",
"title": "首页内存告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.system_push",
"name": "PanelSysMEMTask"
},
"template": {
"field": [
{
"attr": "cycle",
"name": "每",
"type": "select",
"unit": "分钟",
"suffix": "内平均",
"width": "70px",
"disabled": true,
"default": 5,
"items": [
{
"title": "1",
"value": 3
},
{
"title": "5",
"value": 5
},
{
"title": "15",
"value": 15
}
]
},
{
"attr": "count",
"name": "内存使用率超过",
"type": "number",
"unit": "%",
"suffix": "后触发告警",
"default": 80,
"err_msg_prefix": "内存"
}
],
"sorted": [
[
"cycle",
"count"
]
]
},
"default": {
"cycle": 5,
"count": 80
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook",
"sms"
],
"unique": true
}
]
+503
View File
@@ -0,0 +1,503 @@
import json
import os
import sys
import threading
import time
from datetime import datetime, timedelta
from importlib import import_module
from typing import Tuple, Union, Optional, List
import psutil
from .send_tool import WxAccountMsg
from .base_task import BaseTask
from .mods import PUSH_DATA_PATH, TaskTemplateConfig
from .util import read_file, write_file, get_config_value, GET_CLASS
class _ProcessInfo:
def __init__(self):
self.data = None
self.last_time = 0
def __call__(self) -> list:
if self.data is not None and time.time() - self.last_time < 60:
return self.data
try:
import PluginLoader
get_obj = GET_CLASS()
get_obj.sort = "status"
p_info = PluginLoader.plugin_run("task_manager", "get_process_list", get_obj)
except:
return []
if isinstance(p_info, dict) and "process_list" in p_info and isinstance(
p_info["process_list"], list):
self._process_info = p_info["process_list"]
self.last_time = time.time()
return self._process_info
else:
return []
get_process_info = _ProcessInfo()
def have_task_manager_plugin():
"""
通过文件判断是否有进程管理器
"""
return os.path.exists("/www/server/panel/plugin/task_manager/task_manager_push.py")
def load_task_manager_template():
if TaskTemplateConfig().get_by_id("60"):
return None
from .mods import load_task_template_by_config
load_task_template_by_config([
{
"id": "60",
"ver": "1",
"used": True,
"source": "task_manager_cpu",
"title": "任务管理器CPU占用量告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.task_manager_push",
"name": "TaskManagerCPUTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "进程名称",
"type": "select",
"items": {
"url": "plugin?action=a&name=task_manager&s=get_process_list_to_push"
}
},
{
"attr": "count",
"name": "占用率超过",
"type": "number",
"unit": "%",
"suffix": "后触发告警",
"default": 80,
"err_msg_prefix": "CUP占用率"
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"project"
],
[
"count"
],
[
"interval"
]
],
},
"default": {
"project": '',
"count": 80,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": False
},
{
"id": "61",
"ver": "1",
"used": True,
"source": "task_manager_mem",
"title": "任务管理器内存占用量告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.task_manager_push",
"name": "TaskManagerMEMTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "进程名称",
"type": "select",
"items": {
"url": "plugin?action=a&name=task_manager&s=get_process_list_to_push"
}
},
{
"attr": "count",
"name": "占用量超过",
"type": "number",
"unit": "MB",
"suffix": "后触发告警",
"default": None,
"err_msg_prefix": "占用量"
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"project"
],
[
"count"
],
[
"interval"
]
],
},
"default": {
"project": '',
"count": 80,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": False
},
{
"id": "62",
"ver": "1",
"used": True,
"source": "task_manager_process",
"title": "任务管理器进程开销告警",
"load_cls": {
"load_type": "path",
"cls_path": "mod.base.push_mod.task_manager_push",
"name": "TaskManagerProcessTask"
},
"template": {
"field": [
{
"attr": "project",
"name": "进程名称",
"type": "select",
"items": {
"url": "plugin?action=a&name=task_manager&s=get_process_list_to_push"
}
},
{
"attr": "count",
"name": "进程数超过",
"type": "number",
"unit": "个",
"suffix": "后触发告警",
"default": 20,
"err_msg_prefix": "进程数"
},
{
"attr": "interval",
"name": "间隔时间",
"type": "number",
"unit": "秒",
"suffix": "后再次监控检测条件",
"default": 600
}
],
"sorted": [
[
"project"
],
[
"count"
],
[
"interval"
]
],
},
"default": {
"project": '',
"count": 80,
"interval": 600
},
"advanced_default": {
"number_rule": {
"day_num": 3
}
},
"send_type_list": [
"wx_account",
"dingding",
"feishu",
"mail",
"weixin",
"webhook"
],
"unique": False
}
])
class TaskManagerCPUTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "task_manager_cpu"
self.template_name = "任务管理器CUP占用量告警"
def get_title(self, task_data: dict) -> str:
return "进程【{}】的CPU占用量告警".format(task_data["project"])
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if "interval" not in task_data or not isinstance(task_data["interval"], int):
task_data["interval"] = 600
if task_data["interval"] < 60:
task_data["interval"] = 60
if "count" not in task_data or not isinstance(task_data["count"], int):
return "设置的检查范围不正确"
if not 1 <= task_data["count"] < 100:
return "设置的检查范围不正确"
if not task_data["project"]:
return "请选择进程"
return task_data
def get_keyword(self, task_data: dict) -> str:
return task_data["project"]
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
process_info = get_process_info()
self.title = self.get_title(task_data)
count = used = 0
for p in process_info:
if p["name"] == task_data['project']:
used += p["cpu_percent"]
count += 1 if "children" not in p else len(p["children"]) + 1
if used <= task_data['count']:
return None
return {
'msg_list':
[
">通知类型:任务管理器CPU占用量告警",
">告警内容: 进程名称为【{}】的进程共有{}个,消耗的CPU资源占比为{}%,大于告警阈值{}%。".format(
task_data['project'], count, used, task_data['count']
)
],
"project": task_data['project'],
"count": int(task_data['count'])
}
def filter_template(self, template: dict) -> Optional[dict]:
if not have_task_manager_plugin():
return None
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "任务管理器CPU占用量告警"
if len(push_data["project"]) > 11:
project = push_data["project"][:9] + ".."
else:
project = push_data["project"]
msg.msg = "{}的CUP超过{}%".format(project, push_data["count"])
return msg
class TaskManagerMEMTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "task_manager_mem"
self.template_name = "任务管理器内存占用量告警"
def get_title(self, task_data: dict) -> str:
return "进程【{}】的内存占用量告警".format(task_data["project"])
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not task_data["project"]:
return "请选择进程"
if "interval" not in task_data or not isinstance(task_data["interval"], int):
task_data["interval"] = 600
task_data["interval"] = max(60, task_data["interval"])
if "count" not in task_data or not isinstance(task_data["count"], int):
return "设置的检查范围不正确"
if task_data["count"] < 1:
return "设置的检查范围不正确"
return task_data
def get_keyword(self, task_data: dict) -> str:
return task_data["project"]
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
process_info = get_process_info()
self.title = self.get_title(task_data)
used = count = 0
for p in process_info:
if p["name"] == task_data['project']:
used += p["memory_used"]
count += 1 if "children" not in p else len(p["children"]) + 1
if used <= task_data['count'] * 1024 * 1024:
return None
return {
'msg_list': [
">通知类型:任务管理器内存占用量告警",
">告警内容: 进程名称为【{}】的进程共有{}个,消耗的内存资源为{}MB,大于告警阈值{}MB。".format(
task_data['project'], count, int(used / 1024 / 1024), task_data['count']
)
],
"project": task_data['project']
}
def filter_template(self, template: dict) -> Optional[dict]:
if not have_task_manager_plugin():
return None
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
if len(push_data["project"]) > 11:
project = push_data["project"][:9] + ".."
else:
project = push_data["project"]
msg.thing_type = "任务管理器内存占用量告警"
msg.msg = "{}的内存超过告警数值".format(project)
return msg
class TaskManagerProcessTask(BaseTask):
def __init__(self):
super().__init__()
self.source_name = "task_manager_process"
self.title = "任务管理器进程开销告警"
def get_title(self, task_data: dict) -> str:
return "进程【{}】的子进程开销告警".format(task_data["project"])
def check_task_data(self, task_data: dict) -> Union[dict, str]:
if not task_data["project"]:
return "请选择进程"
if "interval" not in task_data or not isinstance(task_data["interval"], int):
task_data["interval"] = 600
task_data["interval"] = max(60, task_data["interval"])
if "count" not in task_data or not isinstance(task_data["count"], int):
return "设置的检查范围不正确"
if task_data["count"] < 1:
return "设置的检查范围不正确"
return task_data
def get_keyword(self, task_data: dict) -> str:
return task_data["project"]
def get_push_data(self, task_id: str, task_data: dict) -> Optional[dict]:
process_info = get_process_info()
count = 0
for p in process_info:
if p["name"] == task_data['project']:
count += 1 if "children" not in p else len(p["children"]) + 1
if count <= task_data['count']:
return None
return {
'msg_list':
[
">通知类型:任务管理器进程开销告警",
">告警内容: 进程名称为【{}】的进程共有{}个,大于告警阈值{}个。".format(
task_data['project'], count, task_data['count']
)
],
"project": task_data['project'],
"count": task_data['count'],
}
def filter_template(self, template: dict) -> Optional[dict]:
if not have_task_manager_plugin():
return None
return template
def to_sms_msg(self, push_data: dict, push_public_data: dict) -> Tuple[str, dict]:
return '', {}
def to_wx_account_msg(self, push_data: dict, push_public_data: dict) -> WxAccountMsg:
msg = WxAccountMsg.new_msg()
msg.thing_type = "任务管理器进程开销告警"
if len(push_data["project"]) > 11:
project = push_data["project"][:9] + ".."
else:
project = push_data["project"]
if push_data["count"] > 100: # 节省字数
push_data["count"] = "限制"
msg.msg = "{}的子进程数超过{}".format(project, push_data["count"])
return msg
class ViewMsgFormat(object):
_FORMAT = {
"60": (
lambda x: "<span>进程:{}的CUP占用超过{}%触发</span>".format(
x.get("project"), x.get("count")
)
),
"61": (
lambda x: "<span>进程:{}的内存使用率超过{}MB后触发</span>".format(
x.get("project"), x.get("count")
)
),
"62": (
lambda x: "<span>进程:{}的子进程数量超过{}后触发</span>".format(
x.get("project"), x.get("count")
)
),
}
def get_msg(self, task: dict) -> Optional[str]:
if task["template_id"] in self._FORMAT:
return self._FORMAT[task["template_id"]](task["task_data"])
return None
+75
View File
@@ -0,0 +1,75 @@
import sys
from typing import Optional, Type, TypeVar
import traceback
from importlib import import_module
from .base_task import BaseTask
from .util import GET_CLASS, get_client_ip, debug_log
T_CLS = TypeVar('T_CLS', bound=BaseTask)
def load_task_cls_by_function(
name: str,
func_name: str,
is_model: bool = False,
model_index: str = '',
args: Optional[dict] = None,
sub_name: Optional[str] = None,
) -> Optional[Type[T_CLS]]:
"""
从执行函数的结果中获取任务类
@param model_index: 模块来源,例如:新场景就是mod
@param name: 名称
@param func_name: 函数名称
@param is_model: 是否在Model中,不在Model中,就应该在插件中
@param args: 请求这个接口的参数, 默认为空
@param sub_name: 自分类名称, 如果有,则会和主名称name做拼接
@return: 返回None 或者有效的任务类
"""
import PluginLoader
real_name = name
if isinstance(sub_name, str):
real_name = "{}/{}".format(name, sub_name)
get_obj = GET_CLASS()
if args is not None and isinstance(args, dict):
for key, value in args.items():
setattr(get_obj, key, value)
try:
if is_model:
get_obj.model_index = model_index
res = PluginLoader.module_run(real_name, func_name, get_obj)
else:
get_obj.fun = func_name
get_obj.s = func_name
get_obj.client_ip = get_client_ip
res = PluginLoader.plugin_run(name, func_name, get_obj)
except:
debug_log(traceback.format_exc())
return None
if isinstance(res, dict):
return None
elif isinstance(res, BaseTask):
return res.__class__
elif issubclass(res, BaseTask):
return res
return None
def load_task_cls_by_path(path: str, cls_name: str) -> Optional[Type[T_CLS]]:
try:
module = import_module(path)
cls = getattr(module, cls_name, None)
if issubclass(cls, BaseTask):
return cls
elif isinstance(cls, BaseTask):
return cls.__class__
else:
return None
except:
print(traceback.format_exc())
print(sys.path)
debug_log(traceback.format_exc())
return None
+98
View File
@@ -0,0 +1,98 @@
import sys
import time
from typing import Optional, Callable
if "/www/server/panel/class" not in sys.path:
sys.path.insert(0, "/www/server/panel/class")
import public
from db import Sql
def write_file(filename: str, s_body: str, mode='w+') -> bool:
"""
写入文件内容
@filename 文件名
@s_body 欲写入的内容
return bool 若文件不存在则尝试自动创建
"""
try:
fp = open(filename, mode=mode)
fp.write(s_body)
fp.close()
return True
except:
try:
fp = open(filename, mode=mode, encoding="utf-8")
fp.write(s_body)
fp.close()
return True
except:
return False
def read_file(filename, mode='r') -> Optional[str]:
"""
读取文件内容
@filename 文件名
return string(bin) 若文件不存在,则返回None
"""
import os
if not os.path.exists(filename):
return None
fp = None
try:
fp = open(filename, mode=mode)
f_body = fp.read()
except:
return None
finally:
if fp and not fp.closed:
fp.close()
return f_body
ExecShell: Callable = public.ExecShell
write_log: Callable = public.WriteLog
Sqlite: Callable = Sql
GET_CLASS: Callable = public.dict_obj
debug_log: Callable = public.print_log
get_config_value: Callable = public.GetConfigValue
get_server_ip: Callable = public.get_server_ip
get_network_ip: Callable = public.get_network_ip
format_date: Callable = public.format_date
public_get_cache_func: Callable = public.get_cache_func
public_set_cache_func: Callable = public.set_cache_func
public_get_user_info: Callable = public.get_user_info
public_http_post = public.httpPost
panel_version = public.version
def get_client_ip() -> str:
return public.GetClientIp()
class _DB:
def __call__(self, table: str):
import db
with db.Sql() as t:
t.table(table)
return t
DB = _DB()