Files
GX-gp-notify/gx_gp_monitor/main.py
T

505 lines
19 KiB
Python

#!/usr/bin/env python3
"""
广西政府采购网公告监控系统主程序
广西政府采购网公告爬取和监控的智能系统
"""
import sys
import argparse
import signal
from pathlib import Path
from typing import Dict, Any
# 添加项目根目录到Python路径
project_root = Path(__file__).parent
sys.path.insert(0, str(project_root))
try:
# 尝试相对导入
from .core.config_manager import load_config, get_config
from .core.logger import init_logger, get_logger
from .core.reliability import check_system_health
from .crawler.spider import crawl_announcements
from .filters.filters import filter_from_config
from .storage.postgresql import init_storage, save_announcements_to_storage, save_all_announcements_by_source_to_storage, cleanup_storage
from .storage.md_generator import generate_onu_md
from .notification.wechat import send_announcements_notification, send_system_notification
from .wechat.callback_server import get_callback_server
from .wechat.menu_manager import WeChatMenuManager
except ImportError:
try:
# 尝试绝对导入(直接运行脚本时)
from core.config_manager import load_config, get_config
from core.logger import init_logger, get_logger
from core.reliability import check_system_health
from crawler.spider import crawl_announcements
from filters.filters import filter_from_config
from storage.postgresql import init_storage, save_announcements_to_storage, save_all_announcements_by_source_to_storage, cleanup_storage
from storage.md_generator import generate_onu_md
from notification.wechat import send_announcements_notification, send_system_notification
# 企业微信模块动态导入,避免循环导入问题
wechat_available = True
try:
from wechat.callback_server import get_callback_server
from wechat.menu_manager import WeChatMenuManager
except ImportError:
wechat_available = False
print("企业微信模块不可用", file=sys.stderr)
except ImportError as e:
print(f"导入错误: {e}", file=sys.stderr)
print("请确保依赖已正确安装: pip install -r requirements.txt", file=sys.stderr)
sys.exit(1)
logger = get_logger(__name__)
class GXGPMonitorApp:
"""广西政府采购网监控系统应用"""
def __init__(self):
self.config = None
self.logger = None
self.running = False
def initialize(self, config_file: str = None):
"""初始化应用"""
try:
# 加载配置
self.config = load_config(config_file)
# 初始化日志
self.logger = init_logger(config=self.config)
logger.info("=== 广西政府采购网公告监控系统启动 ===")
logger.info(f"版本: 1.0.0")
logger.info(f"配置文件: {config_file or '默认配置'}")
# 系统健康检查
if not check_system_health():
logger.warning("系统健康检查失败,但继续运行")
# 初始化存储
init_storage()
logger.info("应用初始化完成")
return True
except Exception as e:
print(f"应用初始化失败: {str(e)}", file=sys.stderr)
return False
def run_crawl(self, keywords: list = None, sources: list = None, max_pages: int = None, manual_crawl: bool = False):
"""执行爬取任务
Args:
keywords: 关键词列表
sources: 来源列表
max_pages: 最大页数
manual_crawl: 是否为手动爬取(只筛选今天的数据)
"""
try:
logger.info(f"开始执行爬取任务 (手动爬取: {manual_crawl})")
# 执行爬取
crawl_results = crawl_announcements(keywords, sources)
if not crawl_results:
logger.info("爬取完成:无数据")
return {"success": True, "results": []}
# 收集所有公告
all_announcements = []
for result in crawl_results:
if result.announcements:
all_announcements.extend(result.announcements)
logger.info(f"爬取到 {len(all_announcements)} 条原始公告")
# 对于手动爬取,不保存公告到数据库,只进行筛选和返回结果
if not manual_crawl:
# 先保存所有公告(按来源分组,每源保留最新100条)
all_saved_stats = save_all_announcements_by_source_to_storage(all_announcements, max_per_source=100)
all_saved_count = sum(all_saved_stats.values())
logger.info(f"保存所有公告完成:共保存 {all_saved_count} 条,按来源统计: {all_saved_stats}")
else:
all_saved_count = 0
all_saved_stats = {}
logger.info("手动爬取模式:跳过数据库保存")
# 筛选公告
if manual_crawl:
# 手动爬取:只筛选今天的公告和用户指定的关键词,不进行去重
from .filters.filters import KeywordFilter, DateFilter, SourceFilter
from datetime import date
# 1. 日期筛选:只保留今天的公告
date_filter = DateFilter()
date_filtered = date_filter.filter_announcements(all_announcements, start_date=date.today(), end_date=date.today())
# 2. 关键词筛选
keyword_filter = KeywordFilter()
keyword_filtered = keyword_filter.filter_announcements(date_filtered, keywords=keywords or [])
# 3. 来源筛选
source_filter = SourceFilter()
filtered_announcements = source_filter.filter_announcements(keyword_filtered, sources or list(self.config.sources.keys()))
# 计算统计信息
filter_stats = type('FilterResult', (), {
"keyword_filtered": len(date_filtered) - len(keyword_filtered),
"date_filtered": len(all_announcements) - len(date_filtered),
"duplicate_filtered": 0, # 手动爬取不进行去重
"source_filtered": len(keyword_filtered) - len(filtered_announcements)
})()
else:
# 自动爬取:筛选出新公告并应用关键词筛选
# 1. 筛选出数据库中没有的新公告
new_announcements = [ann for ann in all_announcements if ann.is_new]
logger.info(f"筛选出 {len(new_announcements)} 条新公告")
# 2. 对新公告应用关键词筛选等
if new_announcements:
filter_obj = filter_from_config()
filtered_announcements, filter_stats = filter_obj.filter(new_announcements)
# 更新统计信息,加上未筛选的新公告数量
filter_stats.keyword_filtered += len(new_announcements) - len(filtered_announcements)
else:
filtered_announcements = []
filter_stats = type('FilterResult', (), {
"keyword_filtered": 0,
"date_filtered": 0,
"duplicate_filtered": len(all_announcements),
"source_filtered": 0
})()
logger.info(f"筛选后剩余 {len(filtered_announcements)} 条公告")
# 对于手动爬取,不保存筛选后的公告到数据库
if not manual_crawl:
# 保存筛选后的公告(用于标记关键词匹配等)
saved_count = save_announcements_to_storage(filtered_announcements)
else:
saved_count = 0
logger.info("手动爬取模式:跳过筛选后公告的数据库保存")
# 生成Markdown文件(只在自动爬取时生成)
md_success = False
if not manual_crawl:
md_success = generate_onu_md(filtered_announcements)
# 发送通知
notify_success = False
if not manual_crawl and filtered_announcements and self.config.wechat_app.enabled:
# 自动爬取时发送卡片消息
notify_success = send_announcements_notification(filtered_announcements)
elif manual_crawl and filtered_announcements and self.config.wechat_app.enabled:
# 手动爬取时不在这里发送消息,由消息处理器负责发送markdown消息
notify_success = True # 标记为成功,因为消息会通过其他方式发送
result = {
"success": True,
"total_crawled": len(all_announcements),
"filtered": len(filtered_announcements),
"saved": saved_count,
"markdown_generated": md_success,
"notification_sent": notify_success,
"filter_stats": {
"keyword_filtered": filter_stats.keyword_filtered,
"date_filtered": filter_stats.date_filtered,
"duplicate_filtered": filter_stats.duplicate_filtered,
"source_filtered": filter_stats.source_filtered
}
}
# 对于手动爬取,额外返回筛选后的公告列表
if manual_crawl:
result["filtered_announcements"] = filtered_announcements
logger.info(f"爬取任务完成: {result}")
return result
except Exception as e:
logger.error(f"爬取任务执行失败: {str(e)}")
return {"success": False, "error": str(e)}
def run_cleanup(self, days: int = None):
"""执行数据清理任务"""
try:
logger.info("开始执行数据清理任务")
deleted_count = cleanup_storage(days)
logger.info(f"数据清理完成:删除 {deleted_count} 条过期数据")
# 发送通知
if deleted_count > 0 and self.config.wechat_app.enabled:
send_system_notification(
"数据清理完成",
f"已清理 {deleted_count} 条过期数据"
)
return {"success": True, "deleted": deleted_count}
except Exception as e:
logger.error(f"数据清理任务执行失败: {str(e)}")
return {"success": False, "error": str(e)}
def run_wechat_server(self, host: str = '0.0.0.0', port: int = 18001):
"""启动企业微信回调服务器"""
if not wechat_available:
return {"success": False, "error": "企业微信模块不可用"}
try:
logger.info("启动企业微信回调服务器")
# 获取回调服务器
callback_server = get_callback_server()
# 启动服务器
callback_server.run(host=host, port=port, debug=self.config.debug)
return {"success": True, "host": host, "port": port}
except Exception as e:
logger.error(f"启动企业微信回调服务器失败: {str(e)}")
return {"success": False, "error": str(e)}
def manage_wechat_menu(self, action: str) -> Dict[str, Any]:
"""管理企业微信菜单"""
if not wechat_available:
return {"success": False, "error": "企业微信模块不可用"}
try:
logger.info(f"执行企业微信菜单操作: {action}")
menu_manager = WeChatMenuManager()
if action == 'create':
success = menu_manager.create_menu()
result = {"success": success, "action": "create"}
elif action == 'delete':
success = menu_manager.delete_menu()
result = {"success": success, "action": "delete"}
elif action == 'get':
menu_info = menu_manager.get_menu()
result = {"success": menu_info is not None, "action": "get", "menu": menu_info}
elif action == 'test':
test_results = menu_manager.test_menu_operations()
result = {"success": True, "action": "test", "results": test_results}
else:
result = {"success": False, "error": f"未知操作: {action}"}
if result["success"]:
logger.info(f"企业微信菜单操作成功: {action}")
else:
logger.error(f"企业微信菜单操作失败: {action}")
return result
except Exception as e:
logger.error(f"企业微信菜单管理异常: {str(e)}")
return {"success": False, "error": str(e)}
def show_status(self):
"""显示系统状态"""
try:
status = {
"system": {
"version": "1.0.0",
"healthy": check_system_health()
},
"config": {
"debug": self.config.debug,
"log_level": self.config.log_level.value
},
"database": {
"enabled": self.config.database.enabled,
"type": self.config.database.type
},
"wechat": {
"enabled": self.config.wechat_app.enabled
}
}
# 格式化输出
print("\n=== 系统状态 ===")
print(f"系统健康: {'正常' if status['system']['healthy'] else '异常'}")
print(f"调试模式: {'开启' if status['config']['debug'] else '关闭'}")
print(f"日志级别: {status['config']['log_level']}")
print(f"数据库: {'启用' if status['database']['enabled'] else '禁用'} ({status['database']['type']})")
print(f"企业微信: {'启用' if status['wechat']['enabled'] else '禁用'}")
return status
except Exception as e:
logger.error(f"获取系统状态失败: {str(e)}")
return {"error": str(e)}
def create_argument_parser():
"""创建命令行参数解析器"""
parser = argparse.ArgumentParser(
description="广西政府采购网公告监控系统",
formatter_class=argparse.RawDescriptionHelpFormatter,
epilog="""
使用示例:
python main.py crawl # 执行一次爬取
python main.py crawl --keywords "大化" # 爬取指定关键词
python main.py cleanup # 执行数据清理
python main.py status # 查看系统状态
python main.py wechat-server # 启动企业微信回调服务器
python main.py wechat-menu --action create # 创建企业微信菜单
"""
)
parser.add_argument(
'command',
choices=['crawl', 'cleanup', 'status', 'wechat-server', 'wechat-menu'],
help='要执行的命令'
)
parser.add_argument(
'--config', '-c',
help='配置文件路径'
)
# crawl命令的参数
parser.add_argument(
'--keywords', '-k',
nargs='+',
help='关键词过滤(多个关键词用空格分隔)'
)
parser.add_argument(
'--sources', '-s',
nargs='+',
help='来源代码过滤(多个来源用空格分隔)'
)
parser.add_argument(
'--max-pages',
type=int,
help='最大爬取页数'
)
# cleanup命令的参数
parser.add_argument(
'--days', '-d',
type=int,
help='清理多少天前的过期数据'
)
# wechat-server命令的参数
parser.add_argument(
'--host', '-H',
default='0.0.0.0',
help='服务器监听主机地址 (默认: 0.0.0.0)'
)
parser.add_argument(
'--port', '-P',
type=int,
default=18001,
help='服务器监听端口 (默认: 18001)'
)
# wechat-menu命令的参数
parser.add_argument(
'--action', '-a',
choices=['create', 'delete', 'get', 'test'],
default='create',
help='菜单操作类型 (默认: create)'
)
return parser
def main():
"""主函数"""
parser = create_argument_parser()
args = parser.parse_args()
# 创建应用实例
app = GXGPMonitorApp()
# 初始化应用
if not app.initialize(args.config):
sys.exit(1)
try:
if args.command == 'crawl':
# 执行爬取
result = app.run_crawl(
keywords=args.keywords,
sources=args.sources,
max_pages=args.max_pages
)
if result["success"]:
print("✅ 爬取任务执行成功")
print(f" 爬取公告: {result['total_crawled']}")
print(f" 筛选后: {result['filtered']}")
print(f" 保存数量: {result['saved']}")
if result.get("markdown_generated"):
print(" Markdown文件: 已生成")
if result.get("notification_sent"):
print(" 通知发送: 已发送")
else:
print(f"❌ 爬取任务执行失败: {result.get('error', '未知错误')}")
sys.exit(1)
elif args.command == 'cleanup':
# 执行清理
result = app.run_cleanup(days=args.days)
if result["success"]:
print(f"✅ 数据清理完成,删除 {result['deleted']} 条记录")
else:
print(f"❌ 数据清理失败: {result.get('error', '未知错误')}")
sys.exit(1)
elif args.command == 'status':
# 显示状态
app.show_status()
elif args.command == 'wechat-server':
# 启动企业微信回调服务器
result = app.run_wechat_server(host=args.host, port=args.port)
if result["success"]:
print(f"✅ 企业微信回调服务器已启动: {result['host']}:{result['port']}")
print(" 回调地址: /api/v1/wechat/callback")
else:
print(f"❌ 企业微信回调服务器启动失败: {result.get('error', '未知错误')}")
sys.exit(1)
elif args.command == 'wechat-menu':
# 企业微信菜单管理
result = app.manage_wechat_menu(action=args.action)
if result["success"]:
if args.action == 'create':
print("✅ 企业微信菜单创建成功")
elif args.action == 'delete':
print("✅ 企业微信菜单删除成功")
elif args.action == 'get':
print("✅ 企业微信菜单获取成功")
if result.get("menu"):
print("菜单信息:", json.dumps(result["menu"], indent=2, ensure_ascii=False))
elif args.action == 'test':
print("✅ 企业微信菜单测试完成")
print("测试结果:", result.get("results"))
else:
error = result.get("error", "未知错误")
print(f"❌ 企业微信菜单操作失败: {error}")
sys.exit(1)
except KeyboardInterrupt:
logger.info("收到中断信号,正在退出...")
except Exception as e:
logger.error(f"程序执行异常: {str(e)}")
print(f"❌ 程序执行异常: {str(e)}", file=sys.stderr)
sys.exit(1)
if __name__ == "__main__":
main()