feat(wechat): 新增企微菜单功能并修复重复触发问题
- 新增「查询」菜单组:监控配置、系统状态、今日工作日、最新公告 - 新增「系统管理」子菜单:暂停/恢复定时任务处理器 - 修复立即爬取被企微重试机制触发多次的问题(60秒防重入锁) - docker-compose 补充 image 名称和 container_name Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -1,4 +1,6 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
import xml.etree.ElementTree as ET
|
||||
|
||||
from app.config import settings
|
||||
@@ -7,6 +9,10 @@ from app.wechat.client import WeChatClient
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 防重入:记录最近一次触发爬取的时间戳,60秒内不重复执行
|
||||
_last_crawl_time: float = 0.0
|
||||
_crawl_lock = asyncio.Lock()
|
||||
|
||||
|
||||
class WeChatMessageHandler:
|
||||
def __init__(self):
|
||||
@@ -65,6 +71,18 @@ class WeChatMessageHandler:
|
||||
return await self._handle_trigger_crawl(from_user)
|
||||
elif event_key == "sync_holidays":
|
||||
return await self._handle_sync_holidays(from_user)
|
||||
elif event_key == "monitor_config":
|
||||
return await self._handle_monitor_config(from_user)
|
||||
elif event_key == "system_status":
|
||||
return await self._handle_system_status(from_user)
|
||||
elif event_key == "workday_status":
|
||||
return await self._handle_workday_status(from_user)
|
||||
elif event_key == "latest_announcements":
|
||||
return await self._handle_latest_announcements(from_user)
|
||||
elif event_key == "pause_scheduler":
|
||||
return await self._handle_pause_scheduler(from_user)
|
||||
elif event_key == "resume_scheduler":
|
||||
return await self._handle_resume_scheduler(from_user)
|
||||
return None
|
||||
|
||||
async def handle_text(self, content: str, from_user: str) -> str | None:
|
||||
@@ -97,8 +115,17 @@ class WeChatMessageHandler:
|
||||
await self.client.send_text("查询失败,请稍后再试", from_user)
|
||||
|
||||
async def _handle_trigger_crawl(self, from_user: str) -> str | None:
|
||||
global _last_crawl_time
|
||||
from app.api.deps import get_crawl_service
|
||||
|
||||
# 防重入:企业微信会对同一事件重试多次,60秒内只执行一次
|
||||
async with _crawl_lock:
|
||||
now = time.monotonic()
|
||||
if now - _last_crawl_time < 60:
|
||||
logger.info(f"爬取请求被忽略(防重入),距上次 {now - _last_crawl_time:.1f}s")
|
||||
return None
|
||||
_last_crawl_time = now
|
||||
|
||||
await self.client.send_text("开始爬取,请稍候...", from_user)
|
||||
|
||||
try:
|
||||
@@ -136,3 +163,126 @@ class WeChatMessageHandler:
|
||||
except Exception as e:
|
||||
logger.error(f"同步节假日失败: {e}")
|
||||
await self.client.send_text(f"同步失败: {e}", from_user)
|
||||
|
||||
async def _handle_monitor_config(self, from_user: str) -> str | None:
|
||||
import json
|
||||
try:
|
||||
keywords = settings.crawler_keywords
|
||||
sources = json.loads(settings.announcement_sources)
|
||||
source_names = "、".join(v["name"] for v in sources.values())
|
||||
text = (
|
||||
f"监控关键词: {', '.join(keywords)}\n"
|
||||
f"爬取页数: {settings.crawler_max_pages} 页\n"
|
||||
f"定时规则: {settings.scheduler_cron}\n"
|
||||
f"公告来源: {source_names}"
|
||||
)
|
||||
await self.client.send_text(text, from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"查询监控配置失败: {e}")
|
||||
await self.client.send_text("查询失败,请稍后再试", from_user)
|
||||
|
||||
async def _handle_system_status(self, from_user: str) -> str | None:
|
||||
from app.api.deps import get_db
|
||||
from sqlalchemy import func, select
|
||||
from app.models.announcement import Announcement
|
||||
|
||||
try:
|
||||
async for db in get_db():
|
||||
total_result = await db.execute(
|
||||
select(func.count()).select_from(Announcement)
|
||||
)
|
||||
total = total_result.scalar() or 0
|
||||
|
||||
today_result = await db.execute(
|
||||
select(func.count()).where(
|
||||
func.date(Announcement.publish_date) == func.current_date()
|
||||
).select_from(Announcement)
|
||||
)
|
||||
today = today_result.scalar() or 0
|
||||
|
||||
unsent_result = await db.execute(
|
||||
select(func.count()).where(
|
||||
Announcement.is_sent == False, # noqa: E712
|
||||
Announcement.keyword_matched == True, # noqa: E712
|
||||
).select_from(Announcement)
|
||||
)
|
||||
unsent = unsent_result.scalar() or 0
|
||||
|
||||
scheduler_status = "已启用" if settings.scheduler_enabled else "已禁用"
|
||||
text = (
|
||||
f"累计公告: {total} 条\n"
|
||||
f"今日新增: {today} 条\n"
|
||||
f"待推送: {unsent} 条\n"
|
||||
f"定时任务: {scheduler_status}\n"
|
||||
f"定时规则: {settings.scheduler_cron}"
|
||||
)
|
||||
await self.client.send_text(text, from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"查询系统状态失败: {e}")
|
||||
await self.client.send_text("查询失败,请稍后再试", from_user)
|
||||
|
||||
async def _handle_workday_status(self, from_user: str) -> str | None:
|
||||
from app.services.holiday_service import now_in_china, is_workday
|
||||
from app.api.deps import get_db
|
||||
|
||||
try:
|
||||
async for db in get_db():
|
||||
today = now_in_china()
|
||||
workday = await is_workday(db, today)
|
||||
status = "工作日,正常爬取" if workday else "非工作日,跳过爬取"
|
||||
text = f"今天 {today.strftime('%Y-%m-%d %A')}\n{status}"
|
||||
await self.client.send_text(text, from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"查询工作日状态失败: {e}")
|
||||
await self.client.send_text("查询失败,请稍后再试", from_user)
|
||||
|
||||
async def _handle_latest_announcements(self, from_user: str) -> str | None:
|
||||
from app.api.deps import get_db
|
||||
from sqlalchemy import select, desc
|
||||
from app.models.announcement import Announcement
|
||||
|
||||
try:
|
||||
async for db in get_db():
|
||||
stmt = (
|
||||
select(Announcement)
|
||||
.where(Announcement.keyword_matched == True) # noqa: E712
|
||||
.order_by(desc(Announcement.publish_date))
|
||||
.limit(5)
|
||||
)
|
||||
result = await db.execute(stmt)
|
||||
items = result.scalars().all()
|
||||
|
||||
if not items:
|
||||
await self.client.send_text("暂无匹配关键词的公告", from_user)
|
||||
return None
|
||||
|
||||
lines = ["最新匹配公告(最近5条):"]
|
||||
for i, ann in enumerate(items, 1):
|
||||
date_str = ann.publish_date.strftime("%m-%d") if ann.publish_date else "?"
|
||||
title = ann.title[:30] + "..." if len(ann.title) > 30 else ann.title
|
||||
lines.append(f"{i}. [{date_str}] {title}")
|
||||
await self.client.send_text("\n".join(lines), from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"查询最新公告失败: {e}")
|
||||
await self.client.send_text("查询失败,请稍后再试", from_user)
|
||||
|
||||
async def _handle_pause_scheduler(self, from_user: str) -> str | None:
|
||||
try:
|
||||
from app.scheduler.jobs import scheduler
|
||||
if scheduler.running:
|
||||
scheduler.pause()
|
||||
await self.client.send_text("定时任务已暂停", from_user)
|
||||
else:
|
||||
await self.client.send_text("定时任务未在运行", from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"暂停定时任务失败: {e}")
|
||||
await self.client.send_text(f"操作失败: {e}", from_user)
|
||||
|
||||
async def _handle_resume_scheduler(self, from_user: str) -> str | None:
|
||||
try:
|
||||
from app.scheduler.jobs import scheduler
|
||||
scheduler.resume()
|
||||
await self.client.send_text("定时任务已恢复", from_user)
|
||||
except Exception as e:
|
||||
logger.error(f"恢复定时任务失败: {e}")
|
||||
await self.client.send_text(f"操作失败: {e}", from_user)
|
||||
|
||||
@@ -13,6 +13,31 @@ MENU = {
|
||||
"type": "click",
|
||||
"key": "today_stats",
|
||||
},
|
||||
{
|
||||
"name": "查询",
|
||||
"sub_button": [
|
||||
{
|
||||
"name": "监控配置",
|
||||
"type": "click",
|
||||
"key": "monitor_config",
|
||||
},
|
||||
{
|
||||
"name": "系统状态",
|
||||
"type": "click",
|
||||
"key": "system_status",
|
||||
},
|
||||
{
|
||||
"name": "今日工作日",
|
||||
"type": "click",
|
||||
"key": "workday_status",
|
||||
},
|
||||
{
|
||||
"name": "最新公告",
|
||||
"type": "click",
|
||||
"key": "latest_announcements",
|
||||
},
|
||||
],
|
||||
},
|
||||
{
|
||||
"name": "系统管理",
|
||||
"sub_button": [
|
||||
|
||||
@@ -3,6 +3,8 @@ services:
|
||||
build:
|
||||
context: ..
|
||||
dockerfile: docker/Dockerfile
|
||||
image: gx-gp-notify:latest
|
||||
container_name: gx-gp-notify
|
||||
ports:
|
||||
- "18001:8000"
|
||||
env_file:
|
||||
|
||||
Reference in New Issue
Block a user