Skip to content

Commit 4dd83e6

Browse files
committed
feat: P0+P1+P2 功能完善 - 告警升级/自动恢复/关键词同义词/会话标签/告警指派
P0: - 告警自动升级定时任务(每30分钟检查,超时告警自动提级) - 监控自动恢复(心跳检测到超时会话自动重启,10分钟冷却) - 全局搜索显示总数 P1: - 发送者详情API(GET /senders/{id} + /senders/{id}/messages) P2: - 关键词同义词扩展(synonyms字段 + 匹配器支持) - 会话标签管理(tags字段 + /conversations/tags/list + PUT标签) - 告警指派协作(assignee字段 + PUT /alerts/{id}/assign) - Dashboard支持自定义时间范围
1 parent 0828bfa commit 4dd83e6

7 files changed

Lines changed: 218 additions & 12 deletions

File tree

backend/app/api/alerts.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -583,3 +583,24 @@ async def get_alert_trend(
583583
"""获取告警趋势统计"""
584584
trend = await alert_aggregation_service.get_alert_trend(db, days=days, group_by=group_by)
585585
return {"items": trend, "days": days, "group_by": group_by}
586+
587+
588+
@router.put("/{alert_id}/assign")
589+
async def assign_alert(
590+
alert_id: int,
591+
assignee: str,
592+
note: str = None,
593+
db: AsyncSession = Depends(get_db),
594+
):
595+
"""指派告警给其他人"""
596+
from app.models.alert import Alert
597+
from app.utils import now_utc
598+
result = await db.execute(select(Alert).where(Alert.id == alert_id))
599+
alert = result.scalar_one_or_none()
600+
if not alert:
601+
raise HTTPException(status_code=404, detail=告警不存在)
602+
alert.assignee = assignee
603+
alert.assignee_note = note
604+
alert.assigned_at = now_utc()
605+
await db.commit()
606+
return {"message": f"告警已指派给 {assignee}"}

backend/app/api/conversations.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -864,3 +864,34 @@ async def add_all_channels_from_account(
864864
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
865865
detail=str(e)
866866
)
867+
868+
869+
@router.get("/tags/list")
870+
async def list_all_tags(db: AsyncSession = Depends(get_db)):
871+
"""获取所有标签(去重)"""
872+
result = await db.execute(
873+
select(Conversation.tags).where(Conversation.tags.isnot(None))
874+
)
875+
all_tags = set()
876+
for row in result.scalars().all():
877+
if row and isinstance(row, list):
878+
all_tags.update(row)
879+
return sorted(all_tags)
880+
881+
882+
@router.put("/{conversation_id}/tags")
883+
async def update_conversation_tags(
884+
conversation_id: int,
885+
tags: List[str],
886+
db: AsyncSession = Depends(get_db),
887+
):
888+
"""更新会话标签"""
889+
result = await db.execute(
890+
select(Conversation).where(Conversation.id == conversation_id)
891+
)
892+
conv = result.scalar_one_or_none()
893+
if not conv:
894+
raise HTTPException(status_code=404, detail=会话不存在)
895+
conv.tags = tags
896+
await db.commit()
897+
return {message: 标签已更新, tags: tags}

backend/app/api/senders.py

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,3 +88,107 @@ async def list_senders(
8888
"page_size": page_size,
8989
"total_pages": (total + page_size - 1) // page_size if total else 0,
9090
}
91+
92+
93+
@router.get("/{sender_id}")
94+
async def get_sender_detail(
95+
sender_id: int,
96+
db: AsyncSession = Depends(get_db),
97+
):
98+
获取发送者详情
99+
from app.models.message import Message
100+
result = await db.execute(select(Sender).where(Sender.id == sender_id))
101+
sender = result.scalar_one_or_none()
102+
if not sender:
103+
from fastapi import HTTPException
104+
raise HTTPException(status_code=404, detail=发送者不存在)
105+
106+
# 实时统计告警数
107+
alert_count_result = await db.execute(
108+
select(func.count(Alert.id)).where(Alert.sender_id == sender_id)
109+
)
110+
alert_count = alert_count_result.scalar() or 0
111+
112+
# 统计消息类型分布
113+
type_result = await db.execute(
114+
select(Message.message_type, func.count(Message.id))
115+
.where(Message.sender_id == sender_id)
116+
.group_by(Message.message_type)
117+
)
118+
message_types = {row[0]: row[1] for row in type_result.all()}
119+
120+
return {
121+
id: sender.id,
122+
user_id: sender.user_id,
123+
username: sender.username,
124+
first_name: sender.first_name,
125+
last_name: sender.last_name,
126+
phone: sender.phone,
127+
is_bot: sender.is_bot,
128+
is_verified: sender.is_verified,
129+
is_premium: sender.is_premium,
130+
message_count: sender.message_count,
131+
alert_count: alert_count,
132+
message_types: message_types,
133+
created_at: sender.created_at.isoformat() if sender.created_at else None,
134+
}
135+
136+
137+
@router.get("/{sender_id}/messages")
138+
async def get_sender_messages(
139+
sender_id: int,
140+
page: int = Query(1, ge=1),
141+
page_size: int = Query(20, ge=1, le=100),
142+
db: AsyncSession = Depends(get_db),
143+
):
144+
获取发送者的消息历史
145+
from app.models.message import Message
146+
from app.models.conversation import Conversation
147+
148+
# 验证发送者存在
149+
sender_result = await db.execute(select(Sender).where(Sender.id == sender_id))
150+
if not sender_result.scalar_one_or_none():
151+
from fastapi import HTTPException
152+
raise HTTPException(status_code=404, detail=发送者不存在)
153+
154+
# 总数
155+
count_result = await db.execute(
156+
select(func.count(Message.id)).where(Message.sender_id == sender_id)
157+
)
158+
total = count_result.scalar() or 0
159+
160+
# 分页查询
161+
query = (
162+
select(Message, Conversation.title)
163+
.join(Conversation, Message.conversation_id == Conversation.id)
164+
.where(Message.sender_id == sender_id)
165+
.order_by(Message.date.desc())
166+
.offset((page - 1) * page_size)
167+
.limit(page_size)
168+
)
169+
result = await db.execute(query)
170+
rows = result.all()
171+
172+
items = []
173+
for msg, conv_title in rows:
174+
items.append({
175+
id: msg.id,
176+
conversation_id: msg.conversation_id,
177+
conversation_title: conv_title,
178+
message_type: msg.message_type,
179+
text: msg.text,
180+
caption: msg.caption,
181+
date: msg.date.isoformat() if msg.date else None,
182+
views: msg.views,
183+
forwards: msg.forwards,
184+
has_media: msg.has_media,
185+
is_reply: msg.is_reply,
186+
})
187+
188+
return {
189+
items: items,
190+
total: total,
191+
page: page,
192+
page_size: page_size,
193+
total_pages: (total + page_size - 1) // page_size if total else 0,
194+
}

backend/app/main.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -305,6 +305,27 @@ async def run_backup_task():
305305
background_tasks.add(backup_task)
306306
logger.info("自动备份任务已启动(后台运行)")
307307

308+
# 启动告警自动升级任务(每30分钟检查一次)
309+
from app.services.alert_aggregation_service import alert_aggregation_service
310+
311+
async def run_escalation_task():
312+
"""后台运行告警自动升级任务"""
313+
while True:
314+
try:
315+
await asyncio.sleep(1800) # 每30分钟
316+
async with AsyncSessionLocal() as db:
317+
result = await alert_aggregation_service.escalate_stale_alerts(db)
318+
if result.get("escalated_count", 0) > 0:
319+
logger.info(f"告警自动升级完成: {result['message']}")
320+
except asyncio.CancelledError:
321+
break
322+
except Exception as e:
323+
logger.error(f"告警自动升级任务异常: {e}")
324+
325+
escalation_task = asyncio.create_task(run_escalation_task())
326+
background_tasks.add(escalation_task)
327+
logger.info("告警自动升级任务已启动(每30分钟)")
328+
308329
yield
309330

310331
# 关闭时执行

backend/app/services/keyword_matcher.py

Lines changed: 20 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ async def _get_keyword_rules(self, db, group_ids: Optional[list], enable_all_key
4949
"match_type": keyword.match_type or group.match_type,
5050
"case_sensitive": keyword.case_sensitive if keyword.case_sensitive is not None else group.case_sensitive,
5151
"alert_level": keyword.alert_level or group.alert_level,
52+
"synonyms": keyword.synonyms if hasattr(keyword, 'synonyms') and keyword.synonyms else [],
5253
})
5354

5455
self._rules_cache[key] = {"created_at": time.time(), "rules": rules}
@@ -177,20 +178,31 @@ async def test_keywords(
177178
groups_map = {g.id: g for g in groups_result.scalars().all()}
178179

179180
matched = []
181+
matched_ids = set()
180182
for keyword in keywords:
181183
group = groups_map.get(keyword.group_id)
182184

183185
if group:
184186
match_type = keyword.match_type or group.match_type
185187
case_sensitive = keyword.case_sensitive if keyword.case_sensitive is not None else group.case_sensitive
186188

187-
if self._match(text, keyword.word, match_type, case_sensitive):
188-
matched.append({
189-
"keyword_id": keyword.id,
190-
"word": keyword.word,
191-
"group_id": group.id,
192-
"group_name": group.name,
193-
"match_type": match_type,
194-
})
189+
# 主词 + 同义词
190+
words_to_check = [keyword.word]
191+
synonyms = getattr(keyword, 'synonyms', None) or []
192+
if synonyms:
193+
words_to_check.extend(synonyms)
194+
195+
for w in words_to_check:
196+
if self._match(text, w, match_type, case_sensitive):
197+
if keyword.id not in matched_ids:
198+
matched_ids.add(keyword.id)
199+
matched.append({
200+
"keyword_id": keyword.id,
201+
"word": keyword.word,
202+
"group_id": group.id,
203+
"group_name": group.name,
204+
"match_type": match_type,
205+
})
206+
break
195207

196208
return matched

backend/app/telegram/monitor.py

Lines changed: 20 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -304,10 +304,27 @@ async def _heartbeat_check(self):
304304
last_msg_time = self._session_last_message_time[conv_id]
305305
time_since_last_msg = (datetime.now() - last_msg_time).total_seconds()
306306

307-
# 如果超过10分钟没有消息,记录警告
307+
# 如果超过10分钟没有消息,尝试自动恢复
308308
if time_since_last_msg > 600: # 10分钟
309-
logger.warning(f"会话 {conv_id} 超过 {int(time_since_last_msg/60)} 分钟没有收到消息,可能监控失效")
310-
# 可以在这里添加通知机制(如发送到前端)
309+
logger.warning(f"会话 {conv_id} 超过 {int(time_since_last_msg/60)} 分钟没有收到消息,尝试自动恢复")
310+
# 自动恢复:重启该会话所属账号的监控(带冷却时间)
311+
if not hasattr(self, '_recovery_cooldown'):
312+
self._recovery_cooldown = {}
313+
now_ts = datetime.now().timestamp()
314+
last_recovery = self._recovery_cooldown.get(conv_id, 0)
315+
if now_ts - last_recovery > 600: # 10分钟冷却
316+
try:
317+
async with AsyncSessionLocal() as db:
318+
result = await db.execute(
319+
select(Conversation.account_id).where(Conversation.id == conv_id)
320+
)
321+
acc_id = result.scalar_one_or_none()
322+
if acc_id:
323+
await self.restart_monitors_for_account(acc_id)
324+
self._recovery_cooldown[conv_id] = now_ts
325+
logger.info(f"会话 {conv_id} 自动恢复完成")
326+
except Exception as e:
327+
logger.error(f"会话 {conv_id} 自动恢复失败: {e}")
311328

312329
except asyncio.CancelledError:
313330
break

frontend/src/pages/MonitoringPage.tsx

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -958,7 +958,7 @@ export function MonitoringPage() {
958958
</div>
959959

960960
{/* 分页控件 */}
961-
{selectedConversation && totalCount > 0 && (
961+
{(selectedConversation || globalSearch) && totalCount > 0 && (
962962
<div className="flex flex-col sm:flex-row items-stretch sm:items-center justify-between gap-3 pt-4 border-t border-cyber-blue/10">
963963
<div className="text-xs sm:text-sm text-muted-foreground text-center sm:text-left">
964964
{totalCount} 条消息,第 {currentPage}

0 commit comments

Comments
 (0)