Files
GX-gp-notify/gx_gp_monitor/wechat/message_handler.py
T

445 lines
16 KiB
Python

"""
企业微信消息处理器
处理用户消息和事件,实现菜单功能
"""
import time
import json
from typing import Optional, Dict, Any, List
from datetime import datetime
try:
from ..core.config_manager import get_config
from ..core.logger import get_logger
from ..notification.wechat import send_system_notification
from ..storage.postgresql import save_all_announcements_by_source_to_storage
from ..storage.md_generator import generate_onu_md
from ..core.models import Announcement
except ImportError:
try:
from core.config_manager import get_config
from core.logger import get_logger
from notification.wechat import send_system_notification
from storage.postgresql import save_all_announcements_by_source_to_storage
from storage.md_generator import generate_onu_md
from core.models import Announcement
except ImportError as e:
raise ImportError(f"消息处理器导入失败: {e}")
logger = get_logger(__name__)
class WeChatMessageHandler:
"""企业微信消息处理器"""
def __init__(self):
self.config = get_config()
self.monitor_app = None
# 菜单配置
self.menu_config = {
"crawl": {
"key": "crawl_now",
"name": "立即爬取",
"description": "立即执行一次公告爬取"
},
"today_summary": {
"key": "today_summary",
"name": "今日总结",
"description": "查看今日公告统计"
},
"custom_crawl": {
"key": "custom_crawl",
"name": "自定义爬取",
"description": "输入关键词进行爬取"
}
}
logger.info("企业微信消息处理器初始化完成")
def _get_monitor_app(self):
"""获取监控应用实例"""
if self.monitor_app is None:
# 动态导入避免循环导入
try:
from ..main import GXGPMonitorApp
self.monitor_app = GXGPMonitorApp()
# 初始化但不启动服务器
if not self.monitor_app.initialize():
logger.error("监控应用初始化失败")
return None
except ImportError:
logger.error("无法导入监控应用")
return None
return self.monitor_app
def handle_event(self, event: str, event_key: Optional[str], from_user: str) -> Optional[str]:
"""处理事件消息"""
try:
logger.info(f"处理事件: {event}, key: {event_key}, user: {from_user}")
if event == 'click':
# 菜单点击事件
if event_key == 'crawl_now':
return self._handle_crawl_now(from_user)
elif event_key == 'today_summary':
return self._handle_today_summary(from_user)
elif event_key.startswith('custom_crawl'):
return self._handle_custom_crawl(event_key, from_user)
else:
return self._create_text_response("未知菜单项", from_user)
elif event == 'subscribe':
# 关注事件
welcome_msg = """欢迎关注广西政府采购网公告监控!
我可以帮您:
• 自动监控最新采购公告
• 筛选您关心的关键词信息
• 及时推送重要更新
点击下方菜单开始使用。"""
return self._create_text_response(welcome_msg, from_user)
elif event == 'unsubscribe':
# 取消关注事件
logger.info(f"用户 {from_user} 取消关注")
return None
else:
logger.info(f"未处理的event类型: {event}")
return None
except Exception as e:
logger.error(f"事件处理异常: {str(e)}")
return self._create_text_response("处理失败,请稍后重试", from_user)
def handle_text_message(self, content: str, from_user: str) -> Optional[str]:
"""处理文本消息"""
try:
logger.info(f"处理文本消息: {content}, user: {from_user}")
# 移除前后空格
content = content.strip()
if content == "帮助" or content == "help":
return self._handle_help(from_user)
elif content.startswith("爬取"):
return self._handle_manual_crawl(content, from_user)
elif content.startswith("总结"):
return self._handle_today_summary(from_user)
elif content.startswith("关键词"):
return self._handle_keyword_search(content, from_user)
else:
# 默认当作关键词搜索
return self._handle_keyword_search(f"关键词 {content}", from_user)
except Exception as e:
logger.error(f"文本消息处理异常: {str(e)}")
return self._create_text_response("处理失败,请稍后重试", from_user)
def handle_other_message(self, msg_type: str, from_user: str) -> Optional[str]:
"""处理其他类型的消息"""
try:
logger.info(f"处理其他消息类型: {msg_type}, user: {from_user}")
if msg_type == 'image':
return self._create_text_response("收到图片消息,但我只能处理文本消息", from_user)
elif msg_type == 'voice':
return self._create_text_response("收到语音消息,但我只能处理文本消息", from_user)
else:
return self._create_text_response(f"收到{msg_type}消息,暂不支持此类型", from_user)
except Exception as e:
logger.error(f"其他消息处理异常: {str(e)}")
return self._create_text_response("处理失败,请稍后重试", from_user)
def _handle_crawl_now(self, from_user: str) -> Optional[str]:
"""处理立即爬取菜单"""
try:
logger.info(f"用户 {from_user} 触发立即爬取")
# 获取监控应用
app = self._get_monitor_app()
if not app:
return self._create_text_response("系统初始化失败,请稍后重试", from_user)
# 执行爬取
result = app.run_crawl()
if result.get("success"):
total = result.get("total_crawled", 0)
filtered = result.get("filtered", 0)
saved = result.get("saved", 0)
response = f"""✅ 爬取完成!
📊 统计信息:
• 总共发现: {total} 条公告
• 关键词筛选: {filtered}
• 已保存: {saved}
如有匹配的公告,我会及时推送通知。"""
else:
error = result.get("error", "未知错误")
response = f"❌ 爬取失败: {error}"
return self._create_text_response(response, from_user)
except Exception as e:
logger.error(f"立即爬取处理异常: {str(e)}")
return self._create_text_response("爬取失败,请稍后重试", from_user)
def _handle_today_summary(self, from_user: str) -> Optional[str]:
"""处理今日总结菜单"""
try:
logger.info(f"用户 {from_user} 请求今日总结")
# 这里可以查询今日的公告统计
# 由于数据库查询较为复杂,这里先返回简单的响应
response = """📅 今日公告统计
由于系统正在优化中,今日统计功能暂时不可用。
您可以:
• 点击"立即爬取"获取最新数据
• 发送关键词进行搜索
• 发送"帮助"查看更多功能"""
return self._create_text_response(response, from_user)
except Exception as e:
logger.error(f"今日总结处理异常: {str(e)}")
return self._create_text_response("获取统计失败,请稍后重试", from_user)
def _handle_custom_crawl(self, event_key: str, from_user: str) -> Optional[str]:
"""处理自定义爬取菜单"""
try:
logger.info(f"用户 {from_user} 触发自定义爬取")
response = """🔍 自定义爬取
请回复您想要搜索的关键词,我将为您执行爬取并筛选相关公告。
例如:
• 大化
• 信息化
• 政府采购
发送关键词开始搜索。"""
return self._create_text_response(response, from_user)
except Exception as e:
logger.error(f"自定义爬取处理异常: {str(e)}")
return self._create_text_response("操作失败,请稍后重试", from_user)
def _handle_help(self, from_user: str) -> Optional[str]:
"""处理帮助命令"""
help_text = """🤖 广西政府采购网公告监控助手
📋 菜单功能:
• 立即爬取 - 执行一次公告爬取
• 今日总结 - 查看今日公告统计
• 自定义爬取 - 输入关键词搜索
💬 文本命令:
• 发送关键词 - 搜索相关公告
• "爬取 [关键词]" - 指定关键词爬取
• "总结" - 查看今日统计
• "帮助" - 显示此帮助信息
📢 自动推送:
系统会自动监控最新公告,并推送匹配关键词的信息。
💡 使用提示:
• 关键词支持多个,用空格分隔
• 公告按时间倒序显示
• 点击公告标题可查看详情"""
return self._create_text_response(help_text, from_user)
def _handle_manual_crawl(self, content: str, from_user: str) -> Optional[str]:
"""处理手动爬取命令"""
try:
# 解析关键词
parts = content.split()
if len(parts) < 2:
return self._create_text_response("请指定爬取关键词,例如:爬取 大化", from_user)
keywords = parts[1:]
logger.info(f"用户 {from_user} 手动爬取关键词: {keywords}")
# 获取监控应用
app = self._get_monitor_app()
if not app:
return self._create_text_response("系统初始化失败,请稍后重试", from_user)
# 执行爬取(手动爬取,只筛选今天的公告)
result = app.run_crawl(keywords=keywords, manual_crawl=True)
if result.get("success"):
total = result.get("total_crawled", 0)
filtered = result.get("filtered", 0)
filtered_announcements = result.get("filtered_announcements", [])
# 获取今天的日期范围
from datetime import datetime, date
today = date.today()
time_period = f"{today.strftime('%Y-%m-%d')} 00:00 至 {datetime.now().strftime('%Y-%m-%d %H:%M')}"
if filtered > 0:
# 生成markdown汇总消息并发送
try:
from ..storage.md_generator import MarkdownGenerator
from ..notification.wechat import send_system_notification
# 生成markdown内容
md_generator = MarkdownGenerator()
title = f"手动爬取结果 - 关键词: {' '.join(keywords)}"
markdown_content = md_generator.generate_markdown(filtered_announcements, title, time_period)
# 发送markdown消息
notify_success = send_system_notification(
title="🔍 搜索完成",
content=markdown_content
)
if notify_success:
# 成功发送markdown消息,返回空响应(不发送额外文本消息)
response = ""
else:
response = f"""✅ 爬取完成!
🔍 搜索条件:
• 关键词: {' '.join(keywords)}
• 时间段: {time_period}
📊 统计结果:
• 总共发现: {total} 条公告
• 匹配筛选: {filtered}
⚠️ 公告汇总推送失败,但数据已生成。"""
except Exception as notify_error:
logger.error(f"生成或发送公告汇总失败: {notify_error}")
# 降级处理:手动构建简单的文本响应
announcement_list = []
for i, ann in enumerate(filtered_announcements[:10], 1): # 最多显示10条
announcement_list.append(f"{i}. {ann.title[:50]}...")
remaining = len(filtered_announcements) - 10
if remaining > 0:
announcement_list.append(f"... 还有 {remaining} 条公告")
response = f"""✅ 爬取完成!
🔍 搜索条件:
• 关键词: {' '.join(keywords)}
• 时间段: {time_period}
📊 统计结果:
• 总共发现: {total} 条公告
• 匹配筛选: {filtered}
📋 匹配公告:
{chr(10).join(announcement_list)}
💡 公告详情已保存,可通过其他方式查看。"""
else:
# 没有找到匹配的公告,发送markdown格式的空结果
try:
from ..storage.md_generator import MarkdownGenerator
from ..notification.wechat import send_system_notification
md_generator = MarkdownGenerator()
title = f"手动爬取结果 - 关键词: {' '.join(keywords)}"
markdown_content = md_generator.generate_markdown([], title, time_period)
notify_success = send_system_notification(
title="🔍 搜索完成",
content=markdown_content
)
if notify_success:
response = ""
else:
response = f"""✅ 爬取完成!
🔍 搜索条件:
• 关键词: {' '.join(keywords)}
• 时间段: {time_period}
📊 统计结果:
• 总共发现: {total} 条公告
• 匹配筛选: 0 条
❌ 在指定时间段内没有找到匹配的公告。"""
except Exception as notify_error:
logger.error(f"生成或发送公告汇总失败: {notify_error}")
response = f"""✅ 爬取完成!
🔍 搜索条件:
• 关键词: {' '.join(keywords)}
• 时间段: {time_period}
📊 统计结果:
• 总共发现: {total} 条公告
• 匹配筛选: 0 条
❌ 在指定时间段内没有找到匹配的公告。"""
else:
error = result.get("error", "未知错误")
response = f"❌ 爬取失败: {error}"
return self._create_text_response(response, from_user)
except Exception as e:
logger.error(f"手动爬取处理异常: {str(e)}")
return self._create_text_response("爬取失败,请稍后重试", from_user)
def _handle_keyword_search(self, content: str, from_user: str) -> Optional[str]:
"""处理关键词搜索"""
try:
# 解析关键词
parts = content.split()
keywords = parts[1:] if len(parts) > 1 else parts
if not keywords:
return self._create_text_response("请提供搜索关键词", from_user)
logger.info(f"用户 {from_user} 关键词搜索: {keywords}")
# 这里可以实现关键词搜索逻辑
# 目前先返回提示信息
response = f"""🔍 关键词搜索
搜索关键词: {' '.join(keywords)}
由于系统正在优化中,搜索功能暂时不可用。
您可以:
• 使用"爬取 [关键词]"执行新的爬取
• 点击菜单中的"立即爬取"
• 发送"帮助"查看更多功能"""
return self._create_text_response(response, from_user)
except Exception as e:
logger.error(f"关键词搜索处理异常: {str(e)}")
return self._create_text_response("搜索失败,请稍后重试", from_user)
def _create_text_response(self, content: str, to_user: str) -> str:
"""创建文本消息响应"""
timestamp = str(int(time.time()))
response_xml = f"""<xml>
<ToUserName><![CDATA[{to_user}]]></ToUserName>
<FromUserName><![CDATA[{self.config.wechat_app.corp_id}]]></FromUserName>
<CreateTime>{timestamp}</CreateTime>
<MsgType><![CDATA[text]]></MsgType>
<Content><![CDATA[{content}]]></Content>
</xml>"""
return response_xml