From ee68e3a595ed086ed4e369c4a57e56f77acab10b Mon Sep 17 00:00:00 2001 From: Rizx <157099086+rizx-r@users.noreply.github.com> Date: Wed, 8 Apr 2026 02:25:30 +0800 Subject: [PATCH 1/5] feat(core): unify channel identity and add NapCat adapter --- app/config.py | 2 + app/graph/executors/conversation.py | 12 +- app/graph/executors/delivery.py | 46 ++-- app/graph/executors/memory.py | 9 +- app/graph/graphs/incoming.py | 9 +- app/graph/graphs/memory_update.py | 3 +- app/graph/graphs/proactive.py | 3 +- app/graph/state.py | 24 +- app/graph/subgraphs/context_loading.py | 44 ++-- app/graph/subgraphs/generation.py | 3 +- app/graph/subgraphs/memory_extraction.py | 2 +- app/graph/tools/delivery.py | 11 +- app/main.py | 3 + app/models/admin.py | 38 +++- app/models/database.py | 66 +++++- app/models/user.py | 24 +- app/routers/admin.py | 84 +++++-- app/routers/setup.py | 15 +- app/routers/wecom.py | 3 +- app/schemas/admin.py | 11 + app/services/__init__.py | 8 + app/services/channel_dispatcher.py | 29 +++ app/services/llm_service.py | 6 +- app/services/memory_service.py | 127 ++++++++--- app/services/napcat_service.py | 132 +++++++++++ app/services/proactive_chat_service.py | 276 +++++++---------------- app/services/runtime_config_service.py | 19 ++ requirements.txt | 1 + 28 files changed, 692 insertions(+), 318 deletions(-) create mode 100644 app/services/channel_dispatcher.py create mode 100644 app/services/napcat_service.py diff --git a/app/config.py b/app/config.py index ee61120..a0b6fd5 100644 --- a/app/config.py +++ b/app/config.py @@ -20,6 +20,8 @@ class Settings(BaseSettings): wecom_secret: str = os.getenv("WECOM_SECRET", "") wecom_token: str = os.getenv("WECOM_TOKEN", "") wecom_encoding_aes_key: str = os.getenv("WECOM_ENCODING_AES_KEY", "") + napcat_ws_url: str = os.getenv("NAPCAT_WS_URL", "") + napcat_ws_token: str = os.getenv("NAPCAT_WS_TOKEN", "") # 智谱 API 配置 zhipu_api_key: str = os.getenv("ZHIPU_API_KEY", "") diff --git a/app/graph/executors/conversation.py b/app/graph/executors/conversation.py index 23a3538..2cd91a5 100644 --- a/app/graph/executors/conversation.py +++ b/app/graph/executors/conversation.py @@ -11,7 +11,8 @@ async def save_conversation( *, - user_id: str, + channel: str, + external_user_id: str, user_message: str, agent_message: str, user_emotion: Dict, @@ -19,7 +20,8 @@ async def save_conversation( memories_used: Optional[Dict] = None, ) -> Optional[int]: return await memory_service.save_conversation( - user_id=user_id, + channel=channel, + external_user_id=external_user_id, user_message=user_message, agent_message=agent_message, user_emotion=user_emotion, @@ -30,7 +32,8 @@ async def save_conversation( def schedule_memory_processing( *, - wecom_user_id: str, + channel: str, + external_user_id: str, conversation_id: Optional[int], user_message: str, agent_message: str, @@ -38,7 +41,8 @@ def schedule_memory_processing( agent_emotion: Dict, ) -> None: memory_service.schedule_memory_processing( - wecom_user_id=wecom_user_id, + channel=channel, + external_user_id=external_user_id, conversation_id=conversation_id, user_message=user_message, agent_message=agent_message, diff --git a/app/graph/executors/delivery.py b/app/graph/executors/delivery.py index 0488adb..6547ec9 100644 --- a/app/graph/executors/delivery.py +++ b/app/graph/executors/delivery.py @@ -10,22 +10,23 @@ from app.models.admin import ProactiveChatLog from app.models.database import SessionLocal from app.models.user import Conversation, User -from app.services.wecom_service import wecom_service +from app.services.channel_dispatcher import channel_dispatcher -async def deliver_incoming_reply(*, to_user: str, content: str) -> Dict[str, object]: +async def deliver_incoming_reply(*, channel: str, external_user_id: str, content: str) -> Dict[str, object]: delivery_result: Dict[str, object] = {"attempted": True, "status": "sent"} try: - await wecom_service.send_text_message(to_user, content) + await channel_dispatcher.send_text(channel, external_user_id, content) except Exception as exc: - print(f"发送企业微信消息失败: {exc}") + print(f"Send message failed: {exc}") delivery_result = {"attempted": True, "status": "failed", "error_message": str(exc)} return delivery_result async def deliver_proactive_outreach( *, - target_wecom_user_id: str, + target_channel: str, + target_external_user_id: str, trigger_type: str, window_key: Optional[str], content: str, @@ -35,15 +36,21 @@ async def deliver_proactive_outreach( error_message = None try: - await wecom_service.send_text_message(target_wecom_user_id, content) + await channel_dispatcher.send_text(target_channel, target_external_user_id, content) status = "sent" - _save_proactive_conversation(target_wecom_user_id=target_wecom_user_id, content=content, sent_at=sent_at) + _save_proactive_conversation( + target_channel=target_channel, + target_external_user_id=target_external_user_id, + content=content, + sent_at=sent_at, + ) except Exception as exc: error_message = str(exc) - print(f"主动聊天发送失败: {exc}") + print(f"Proactive delivery failed: {exc}") finally: _save_proactive_log( - target_wecom_user_id=target_wecom_user_id, + target_channel=target_channel, + target_external_user_id=target_external_user_id, trigger_type=trigger_type, window_key=window_key, content=content, @@ -60,10 +67,20 @@ async def deliver_proactive_outreach( } -def _save_proactive_conversation(*, target_wecom_user_id: str, content: str, sent_at: datetime) -> None: +def _save_proactive_conversation( + *, + target_channel: str, + target_external_user_id: str, + content: str, + sent_at: datetime, +) -> None: db = SessionLocal() try: - user = db.query(User).filter(User.wecom_user_id == target_wecom_user_id).first() + user = ( + db.query(User) + .filter(User.channel == target_channel, User.external_user_id == target_external_user_id) + .first() + ) if not user: return @@ -86,7 +103,8 @@ def _save_proactive_conversation(*, target_wecom_user_id: str, content: str, sen def _save_proactive_log( *, - target_wecom_user_id: str, + target_channel: str, + target_external_user_id: str, trigger_type: str, window_key: Optional[str], content: str, @@ -97,7 +115,8 @@ def _save_proactive_log( db = SessionLocal() try: log = ProactiveChatLog( - target_wecom_user_id=target_wecom_user_id, + target_channel=target_channel, + target_external_user_id=target_external_user_id, trigger_type=trigger_type, window_key=window_key, content=content, @@ -109,3 +128,4 @@ def _save_proactive_log( db.commit() finally: db.close() + diff --git a/app/graph/executors/memory.py b/app/graph/executors/memory.py index ecd3d25..3cbd3f2 100644 --- a/app/graph/executors/memory.py +++ b/app/graph/executors/memory.py @@ -11,10 +11,10 @@ from app.services.memory_service import memory_service -async def load_memory_update_context(wecom_user_id: str) -> Tuple[Dict, List[Dict[str, str]]]: +async def load_memory_update_context(channel: str, external_user_id: str) -> Tuple[Dict, List[Dict[str, str]]]: db = SessionLocal() try: - user = memory_service._get_user_by_wecom_id(db, wecom_user_id) + user = memory_service._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return {}, [] @@ -38,7 +38,8 @@ async def load_memory_update_context(wecom_user_id: str) -> Tuple[Dict, List[Dic async def persist_memory_update( *, - wecom_user_id: str, + channel: str, + external_user_id: str, conversation_id: int, user_message: str, agent_message: str, @@ -47,7 +48,7 @@ async def persist_memory_update( ) -> None: db = SessionLocal() try: - user = memory_service._get_user_by_wecom_id(db, wecom_user_id) + user = memory_service._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return diff --git a/app/graph/graphs/incoming.py b/app/graph/graphs/incoming.py index 817f03b..7dd543f 100644 --- a/app/graph/graphs/incoming.py +++ b/app/graph/graphs/incoming.py @@ -16,7 +16,8 @@ async def _save_conversation_node(state: IncomingGraphState) -> IncomingGraphState: conversation_id = await save_conversation( - user_id=state["user_id"], + channel=state["channel"], + external_user_id=state["external_user_id"], user_message=state["user_content"], agent_message=state["agent_response"], user_emotion=state["user_emotion"], @@ -30,7 +31,8 @@ async def _save_conversation_node(state: IncomingGraphState) -> IncomingGraphSta async def _schedule_memory_node(state: IncomingGraphState) -> IncomingGraphState: schedule_memory_processing( - wecom_user_id=state["user_id"], + channel=state["channel"], + external_user_id=state["external_user_id"], conversation_id=state["conversation_id"], user_message=state["user_content"], agent_message=state["agent_response"], @@ -46,7 +48,8 @@ async def _deliver_reply_node(state: IncomingGraphState) -> IncomingGraphState: delivery_result = await message_delivery_tool.ainvoke( { "delivery_kind": "incoming_reply", - "to_user": state["user_id"], + "channel": state["channel"], + "external_user_id": state["external_user_id"], "content": state["agent_response"], } ) diff --git a/app/graph/graphs/memory_update.py b/app/graph/graphs/memory_update.py index 16a6e57..2febef8 100644 --- a/app/graph/graphs/memory_update.py +++ b/app/graph/graphs/memory_update.py @@ -15,7 +15,8 @@ async def _persist_node(state: MemoryUpdateGraphState) -> MemoryUpdateGraphState: await persist_memory_update( - wecom_user_id=state["wecom_user_id"], + channel=state["channel"], + external_user_id=state["external_user_id"], conversation_id=state["conversation_id"], user_message=state["user_message"], agent_message=state["agent_message"], diff --git a/app/graph/graphs/proactive.py b/app/graph/graphs/proactive.py index 70c2179..a3dd6c0 100644 --- a/app/graph/graphs/proactive.py +++ b/app/graph/graphs/proactive.py @@ -26,7 +26,8 @@ async def _deliver_node(state: ProactiveChatGraphState) -> ProactiveChatGraphSta delivery = await message_delivery_tool.ainvoke( { "delivery_kind": "proactive_outreach", - "to_user": state["target_wecom_user_id"], + "channel": state["target_channel"], + "external_user_id": state["target_external_user_id"], "content": state["reply"], "trigger_type": state["trigger_type"], "window_key": state["window_key"], diff --git a/app/graph/state.py b/app/graph/state.py index f1f9cf7..a40590f 100644 --- a/app/graph/state.py +++ b/app/graph/state.py @@ -12,7 +12,8 @@ class IncomingGraphState(TypedDict): - user_id: str + channel: str + external_user_id: str user_content: str persona_config: Dict response_constraints: Dict @@ -34,7 +35,8 @@ class IncomingGraphState(TypedDict): class PreviewGraphState(TypedDict): preview_mode: Literal["prompt", "reply"] user_message: str - wecom_user_id: Optional[str] + channel: Optional[str] + external_user_id: Optional[str] draft_config: Optional[Dict] persona_config: Dict user_memory: Optional[Dict] @@ -52,7 +54,8 @@ class PreviewGraphState(TypedDict): class ProactiveChatGraphState(TypedDict): - target_wecom_user_id: str + target_channel: str + target_external_user_id: str trigger_type: str window_key: Optional[str] send_delivery: bool @@ -70,7 +73,8 @@ class ProactiveChatGraphState(TypedDict): class MemoryUpdateGraphState(TypedDict): - wecom_user_id: str + channel: str + external_user_id: str conversation_id: int user_message: str agent_message: str @@ -95,7 +99,8 @@ def append_tool_trace(state: Dict, name: str) -> List[str]: def build_incoming_initial_state(payload: Dict[str, object]) -> IncomingGraphState: return { - "user_id": str(payload.get("user_id") or ""), + "channel": str(payload.get("channel") or "wecom").strip() or "wecom", + "external_user_id": str(payload.get("external_user_id") or ""), "user_content": str(payload.get("user_content") or ""), "persona_config": {}, "response_constraints": {}, @@ -119,7 +124,8 @@ def build_preview_initial_state(payload: Dict[str, object]) -> PreviewGraphState return { "preview_mode": payload.get("preview_mode", "prompt"), # type: ignore[typeddict-item] "user_message": str(payload.get("user_message") or ""), - "wecom_user_id": str(payload.get("wecom_user_id") or "").strip() or None, + "channel": str(payload.get("channel") or "").strip() or None, + "external_user_id": str(payload.get("external_user_id") or "").strip() or None, "draft_config": payload.get("draft_config") if isinstance(payload.get("draft_config"), dict) else None, "persona_config": {}, "user_memory": None, @@ -139,7 +145,8 @@ def build_preview_initial_state(payload: Dict[str, object]) -> PreviewGraphState def build_proactive_initial_state(payload: Dict[str, object]) -> ProactiveChatGraphState: return { - "target_wecom_user_id": str(payload.get("target_wecom_user_id") or ""), + "target_channel": str(payload.get("target_channel") or "wecom").strip() or "wecom", + "target_external_user_id": str(payload.get("target_external_user_id") or ""), "trigger_type": str(payload.get("trigger_type") or ""), "window_key": str(payload.get("window_key") or "").strip() or None, "send_delivery": bool(payload.get("send_delivery")), @@ -159,7 +166,8 @@ def build_proactive_initial_state(payload: Dict[str, object]) -> ProactiveChatGr def build_memory_initial_state(payload: Dict[str, object]) -> MemoryUpdateGraphState: return { - "wecom_user_id": str(payload.get("wecom_user_id") or ""), + "channel": str(payload.get("channel") or "wecom").strip() or "wecom", + "external_user_id": str(payload.get("external_user_id") or ""), "conversation_id": int(payload.get("conversation_id") or 0), "user_message": str(payload.get("user_message") or ""), "agent_message": str(payload.get("agent_message") or ""), diff --git a/app/graph/subgraphs/context_loading.py b/app/graph/subgraphs/context_loading.py index e04e143..88aa8d3 100644 --- a/app/graph/subgraphs/context_loading.py +++ b/app/graph/subgraphs/context_loading.py @@ -24,14 +24,15 @@ async def _incoming_load_context(state: IncomingGraphState) -> IncomingGraphState: - await memory_service.get_or_create_user(state["user_id"]) + await memory_service.get_or_create_user(state["channel"], state["external_user_id"]) persona_config = persona_service.get_persona_config() response_constraints = get_response_constraints(state["user_content"], persona_config.get("response_preferences")) - context = await memory_service.get_conversation_context(state["user_id"]) - user_memory = await memory_service.get_user_memory(state["user_id"], query_text=state["user_content"]) - recent_agent_replies = await memory_service.get_recent_agent_replies(state["user_id"], limit=3) + context = await memory_service.get_conversation_context(state["channel"], state["external_user_id"]) + user_memory = await memory_service.get_user_memory(state["channel"], state["external_user_id"], query_text=state["user_content"]) + recent_agent_replies = await memory_service.get_recent_agent_replies(state["channel"], state["external_user_id"], limit=3) context_messages = await memory_service.get_recent_messages( - state["user_id"], + state["channel"], + state["external_user_id"], limit=int(response_constraints["context_limit"]), ) return { @@ -52,7 +53,11 @@ async def _incoming_analyze_emotion(state: IncomingGraphState) -> IncomingGraphS print(f"图执行情绪分析失败,回退到默认情绪: {exc}") user_emotion = {"neutral": 1.0} - agent_emotion = await emotion_engine.update_state(state["user_id"], state["user_content"], user_emotion) + agent_emotion = await emotion_engine.update_state( + memory_service.build_user_key(state["channel"], state["external_user_id"]), + state["user_content"], + user_emotion, + ) return { "user_emotion": user_emotion, "agent_emotion": agent_emotion, @@ -61,18 +66,19 @@ async def _incoming_analyze_emotion(state: IncomingGraphState) -> IncomingGraphS async def _preview_load_context(state: PreviewGraphState) -> PreviewGraphState: - wecom_user_id = str(state.get("wecom_user_id") or "").strip() + channel = str(state.get("channel") or "").strip() + external_user_id = str(state.get("external_user_id") or "").strip() user_memory = None context = {} recent_agent_replies = [] context_messages = [] - if wecom_user_id: - user_memory = await memory_service.get_user_memory(wecom_user_id, query_text=state["user_message"]) + if channel and external_user_id: + user_memory = await memory_service.get_user_memory(channel, external_user_id, query_text=state["user_message"]) if user_memory: - context = await memory_service.get_conversation_context(wecom_user_id) - recent_agent_replies = await memory_service.get_recent_agent_replies(wecom_user_id, limit=3) - context_messages = await memory_service.get_recent_messages(wecom_user_id, limit=4) + context = await memory_service.get_conversation_context(channel, external_user_id) + recent_agent_replies = await memory_service.get_recent_agent_replies(channel, external_user_id, limit=3) + context_messages = await memory_service.get_recent_messages(channel, external_user_id, limit=4) persona_config = state["draft_config"] or persona_service.get_persona_config() response_constraints = get_response_constraints(state["user_message"], persona_config.get("response_preferences")) @@ -90,8 +96,10 @@ async def _preview_load_context(state: PreviewGraphState) -> PreviewGraphState: async def _preview_analyze_emotion(state: PreviewGraphState) -> PreviewGraphState: user_emotion = await glm_service.analyze_emotion(state["user_message"]) + channel = str(state.get("channel") or "").strip() + external_user_id = str(state.get("external_user_id") or "").strip() agent_emotion = await emotion_engine.update_state( - state.get("wecom_user_id") or "__preview__", + memory_service.build_user_key(channel or "__preview__", external_user_id or "__preview__"), state["user_message"], user_emotion, ) @@ -105,9 +113,13 @@ async def _preview_analyze_emotion(state: PreviewGraphState) -> PreviewGraphStat async def _proactive_load_context(state: ProactiveChatGraphState) -> ProactiveChatGraphState: persona_config = persona_service.get_persona_config() proactive_config = proactive_chat_service.get_config() - user_memory = await memory_service.get_user_memory(state["target_wecom_user_id"]) - context = await memory_service.get_conversation_context(state["target_wecom_user_id"]) - recent_agent_replies = await memory_service.get_recent_agent_replies(state["target_wecom_user_id"], limit=3) + user_memory = await memory_service.get_user_memory(state["target_channel"], state["target_external_user_id"]) + context = await memory_service.get_conversation_context(state["target_channel"], state["target_external_user_id"]) + recent_agent_replies = await memory_service.get_recent_agent_replies( + state["target_channel"], + state["target_external_user_id"], + limit=3, + ) response_constraints = get_response_constraints( "我想主动开启一段自然微信聊天", persona_config.get("response_preferences"), diff --git a/app/graph/subgraphs/generation.py b/app/graph/subgraphs/generation.py index e531d9e..3540b34 100644 --- a/app/graph/subgraphs/generation.py +++ b/app/graph/subgraphs/generation.py @@ -249,7 +249,8 @@ async def _proactive_generate_reply(state: ProactiveChatGraphState) -> Proactive async def _proactive_finalize_reply(state: ProactiveChatGraphState) -> ProactiveChatGraphState: return { "reply": state["reply"] or proactive_chat_service._build_fallback_message( - state["target_wecom_user_id"], + state["target_channel"], + state["target_external_user_id"], state["trigger_type"], state["user_memory"], ), diff --git a/app/graph/subgraphs/memory_extraction.py b/app/graph/subgraphs/memory_extraction.py index ca4b746..c2502a8 100644 --- a/app/graph/subgraphs/memory_extraction.py +++ b/app/graph/subgraphs/memory_extraction.py @@ -15,7 +15,7 @@ async def _prepare_context(state: MemoryUpdateGraphState) -> MemoryUpdateGraphState: - existing_memory, recent_messages = await load_memory_update_context(state["wecom_user_id"]) + existing_memory, recent_messages = await load_memory_update_context(state["channel"], state["external_user_id"]) return { "existing_memory": existing_memory, "recent_messages": recent_messages, diff --git a/app/graph/tools/delivery.py b/app/graph/tools/delivery.py index 414e8fb..a7e2b76 100644 --- a/app/graph/tools/delivery.py +++ b/app/graph/tools/delivery.py @@ -14,7 +14,8 @@ class MessageDeliveryToolInput(BaseModel): delivery_kind: Literal["incoming_reply", "proactive_outreach"] = Field(..., description="Delivery scenario") - to_user: str = Field(..., description="Target WeCom user id") + channel: str = Field(..., description="Target channel") + external_user_id: str = Field(..., description="Target external user id") content: str = Field(..., description="Message content") trigger_type: Optional[str] = Field(default=None, description="Proactive trigger type") window_key: Optional[str] = Field(default=None, description="Scheduled window key") @@ -23,17 +24,19 @@ class MessageDeliveryToolInput(BaseModel): @tool("message_delivery_tool", args_schema=MessageDeliveryToolInput) async def message_delivery_tool( delivery_kind: Literal["incoming_reply", "proactive_outreach"], - to_user: str, + channel: str, + external_user_id: str, content: str, trigger_type: Optional[str] = None, window_key: Optional[str] = None, ) -> Dict[str, object]: """Deliver a message through the configured executor boundary.""" if delivery_kind == "incoming_reply": - return await deliver_incoming_reply(to_user=to_user, content=content) + return await deliver_incoming_reply(channel=channel, external_user_id=external_user_id, content=content) return await deliver_proactive_outreach( - target_wecom_user_id=to_user, + target_channel=channel, + target_external_user_id=external_user_id, trigger_type=trigger_type or "", window_key=window_key, content=content, diff --git a/app/main.py b/app/main.py index 2177b45..bc1d4e2 100644 --- a/app/main.py +++ b/app/main.py @@ -14,6 +14,7 @@ from app.config import settings from app.models.database import init_db from app.routers import admin, setup, wecom +from app.services.napcat_service import napcat_service from app.services.proactive_chat_service import proactive_chat_service from app.services.tunnel_service import ( is_invalid_autodetected_tunnel_url, @@ -47,6 +48,7 @@ async def lifespan(app: FastAPI): print(f"🚀 恋爱 Agent 启动中...") print(f"📍 服务地址: http://{settings.server_host}:{settings.server_port}") print(f"🔗 企业微信回调地址: {callback_url}") + await napcat_service.start() proactive_scheduler_task = asyncio.create_task(proactive_chat_service.scheduler_loop()) yield @@ -57,6 +59,7 @@ async def lifespan(app: FastAPI): await proactive_scheduler_task except asyncio.CancelledError: pass + await napcat_service.stop() print("👋 恋爱 Agent 关闭中...") diff --git a/app/models/admin.py b/app/models/admin.py index e312127..84f571b 100644 --- a/app/models/admin.py +++ b/app/models/admin.py @@ -5,6 +5,7 @@ from datetime import datetime from sqlalchemy import Boolean, Column, DateTime, Integer, JSON, String, Text +from sqlalchemy.ext.hybrid import hybrid_property from app.models.user import Base @@ -40,7 +41,8 @@ class ProactiveChatConfig(Base): id = Column(Integer, primary_key=True, autoincrement=True) config_key = Column(String(64), unique=True, nullable=False, comment="配置键") enabled = Column(Boolean, nullable=False, default=False, comment="是否启用主动聊天") - target_wecom_user_id = Column(String(64), comment="目标企业微信用户 ID") + target_channel = Column(String(32), nullable=False, default="wecom", comment="目标渠道") + target_external_user_id = Column(String(128), comment="目标渠道用户 ID") scheduled_windows = Column(JSON, nullable=False, default=list, comment="固定时段窗口") inactivity_trigger_hours = Column(Integer, nullable=False, default=6, comment="多久未互动后主动发起") quiet_hours = Column(JSON, nullable=False, default=dict, comment="免打扰时段") @@ -54,6 +56,19 @@ class ProactiveChatConfig(Base): def __repr__(self): return f"" + @hybrid_property + def target_wecom_user_id(self): + return self.target_external_user_id if self.target_channel == "wecom" else None + + @target_wecom_user_id.setter + def target_wecom_user_id(self, value): + self.target_channel = "wecom" + self.target_external_user_id = value + + @target_wecom_user_id.expression + def target_wecom_user_id(cls): # type: ignore[no-redef] + return cls.target_external_user_id + class ProactiveChatLog(Base): """主动聊天发送日志。""" @@ -61,7 +76,8 @@ class ProactiveChatLog(Base): __tablename__ = "proactive_chat_logs" id = Column(Integer, primary_key=True, autoincrement=True) - target_wecom_user_id = Column(String(64), nullable=False, comment="目标企业微信用户 ID") + target_channel = Column(String(32), nullable=False, default="wecom", comment="目标渠道") + target_external_user_id = Column(String(128), nullable=False, comment="目标渠道用户 ID") trigger_type = Column(String(32), nullable=False, comment="触发类型") window_key = Column(String(32), comment="固定窗口 key") content = Column(Text, comment="发送内容") @@ -72,7 +88,23 @@ class ProactiveChatLog(Base): created_at = Column(DateTime, default=datetime.now, comment="创建时间") def __repr__(self): - return f"" + return ( + f"" + ) + + @hybrid_property + def target_wecom_user_id(self): + return self.target_external_user_id if self.target_channel == "wecom" else None + + @target_wecom_user_id.setter + def target_wecom_user_id(self, value): + self.target_channel = "wecom" + self.target_external_user_id = value + + @target_wecom_user_id.expression + def target_wecom_user_id(cls): # type: ignore[no-redef] + return cls.target_external_user_id class RuntimeConfig(Base): diff --git a/app/models/database.py b/app/models/database.py index 4ef9c6a..74ed972 100644 --- a/app/models/database.py +++ b/app/models/database.py @@ -2,7 +2,7 @@ 数据库初始化 """ -from sqlalchemy import create_engine +from sqlalchemy import create_engine, inspect, text from sqlalchemy.orm import sessionmaker from app.config import settings @@ -34,9 +34,73 @@ def init_db(): """初始化数据库,创建所有表""" Base.metadata.create_all(bind=engine) + _run_compat_migrations() print("✅ 数据库表创建完成") +def _run_compat_migrations() -> None: + inspector = inspect(engine) + with engine.begin() as conn: + _ensure_users_channel_columns(inspector, conn) + _ensure_proactive_channel_columns(inspector, conn) + + +def _ensure_users_channel_columns(inspector, conn) -> None: + tables = inspector.get_table_names() + if "users" not in tables: + return + columns = {item["name"] for item in inspector.get_columns("users")} + if "channel" not in columns: + conn.execute(text("ALTER TABLE users ADD COLUMN channel VARCHAR(32) DEFAULT 'wecom' NOT NULL")) + has_legacy_wecom_id = "wecom_user_id" in columns + if "external_user_id" not in columns: + conn.execute(text("ALTER TABLE users ADD COLUMN external_user_id VARCHAR(128)")) + if has_legacy_wecom_id: + conn.execute(text("UPDATE users SET external_user_id = wecom_user_id WHERE external_user_id IS NULL")) + conn.execute(text("UPDATE users SET channel = 'wecom' WHERE channel IS NULL OR channel = ''")) + if has_legacy_wecom_id: + conn.execute( + text("UPDATE users SET external_user_id = wecom_user_id WHERE external_user_id IS NULL OR external_user_id = ''") + ) + + +def _ensure_proactive_channel_columns(inspector, conn) -> None: + tables = inspector.get_table_names() + if "proactive_chat_configs" in tables: + columns = {item["name"] for item in inspector.get_columns("proactive_chat_configs")} + has_legacy_target = "target_wecom_user_id" in columns + if "target_channel" not in columns: + conn.execute(text("ALTER TABLE proactive_chat_configs ADD COLUMN target_channel VARCHAR(32) DEFAULT 'wecom' NOT NULL")) + if "target_external_user_id" not in columns: + conn.execute(text("ALTER TABLE proactive_chat_configs ADD COLUMN target_external_user_id VARCHAR(128)")) + if has_legacy_target: + conn.execute( + text( + "UPDATE proactive_chat_configs " + "SET target_external_user_id = target_wecom_user_id " + "WHERE target_external_user_id IS NULL" + ) + ) + conn.execute(text("UPDATE proactive_chat_configs SET target_channel = 'wecom' WHERE target_channel IS NULL OR target_channel = ''")) + + if "proactive_chat_logs" in tables: + columns = {item["name"] for item in inspector.get_columns("proactive_chat_logs")} + has_legacy_target = "target_wecom_user_id" in columns + if "target_channel" not in columns: + conn.execute(text("ALTER TABLE proactive_chat_logs ADD COLUMN target_channel VARCHAR(32) DEFAULT 'wecom' NOT NULL")) + if "target_external_user_id" not in columns: + conn.execute(text("ALTER TABLE proactive_chat_logs ADD COLUMN target_external_user_id VARCHAR(128)")) + if has_legacy_target: + conn.execute( + text( + "UPDATE proactive_chat_logs " + "SET target_external_user_id = target_wecom_user_id " + "WHERE target_external_user_id IS NULL" + ) + ) + conn.execute(text("UPDATE proactive_chat_logs SET target_channel = 'wecom' WHERE target_channel IS NULL OR target_channel = ''")) + + def get_db(): """获取数据库会话""" db = SessionLocal() diff --git a/app/models/user.py b/app/models/user.py index 2c03e46..700889b 100644 --- a/app/models/user.py +++ b/app/models/user.py @@ -4,8 +4,9 @@ from datetime import datetime -from sqlalchemy import Boolean, Column, DateTime, ForeignKey, Integer, JSON, String, Text +from sqlalchemy import Boolean, Column, DateTime, ForeignKey, Integer, JSON, String, Text, UniqueConstraint from sqlalchemy.ext.declarative import declarative_base +from sqlalchemy.ext.hybrid import hybrid_property from sqlalchemy.orm import relationship Base = declarative_base() @@ -15,9 +16,13 @@ class User(Base): """用户模型""" __tablename__ = "users" + __table_args__ = ( + UniqueConstraint("channel", "external_user_id", name="uq_users_channel_external_user_id"), + ) id = Column(Integer, primary_key=True, autoincrement=True) - wecom_user_id = Column(String(64), unique=True, nullable=False, comment="企业微信用户ID") + channel = Column(String(32), nullable=False, default="wecom", comment="渠道") + external_user_id = Column(String(128), nullable=False, comment="渠道侧用户 ID") nickname = Column(String(100), comment="用户昵称") avatar_url = Column(String(255), comment="头像URL") @@ -47,7 +52,20 @@ class User(Base): memory_items = relationship("MemoryItem", back_populates="user", cascade="all, delete-orphan") def __repr__(self): - return f"" + return f"" + + @hybrid_property + def wecom_user_id(self): + return self.external_user_id if self.channel == "wecom" else None + + @wecom_user_id.setter + def wecom_user_id(self, value): + self.channel = "wecom" + self.external_user_id = value + + @wecom_user_id.expression + def wecom_user_id(cls): # type: ignore[no-redef] + return cls.external_user_id class Conversation(Base): diff --git a/app/routers/admin.py b/app/routers/admin.py index 650d455..cebdd99 100644 --- a/app/routers/admin.py +++ b/app/routers/admin.py @@ -1,15 +1,13 @@ -""" -管理后台 API +""" +Admin backend API. """ from hmac import compare_digest from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response -from app.graph import run_preview_graph from app.config import settings -from app.services.emotion_engine import emotion_engine # 兼容测试 patch -from app.services.llm_service import glm_service # 兼容测试 patch +from app.graph import run_preview_graph from app.schemas.admin import ( AgentPersonaPayload, LoginRequest, @@ -18,6 +16,8 @@ ProactiveChatPayload, UserMemoryPayload, ) +from app.services.emotion_engine import emotion_engine # test patch compatibility +from app.services.llm_service import glm_service # test patch compatibility from app.services.memory_service import memory_service from app.services.persona_service import persona_service from app.services.proactive_chat_service import proactive_chat_service @@ -26,6 +26,18 @@ router = APIRouter(prefix="/admin-api", tags=["管理后台"]) +def _resolve_channel_identity( + channel: str | None, + external_user_id: str | None, + wecom_user_id: str | None, +) -> tuple[str | None, str | None]: + if channel and external_user_id: + return channel, external_user_id + if wecom_user_id: + return "wecom", wecom_user_id + return channel, external_user_id + + def require_admin(request: Request) -> bool: if not request.session.get("is_admin"): raise HTTPException(status_code=401, detail="Admin authentication required") @@ -70,13 +82,18 @@ async def get_proactive_chat(_: bool = Depends(require_admin)): @router.put("/proactive-chat") async def update_proactive_chat(payload: ProactiveChatPayload, _: bool = Depends(require_admin)): - return proactive_chat_service.save_config(payload.model_dump()) + body = payload.model_dump() + if payload.target_wecom_user_id and not payload.target_external_user_id: + body["target_channel"] = "wecom" + body["target_external_user_id"] = payload.target_wecom_user_id + return proactive_chat_service.save_config(body) @router.post("/proactive-chat/preview") async def preview_proactive_chat(payload: ProactiveChatActionRequest, _: bool = Depends(require_admin)): try: - return await proactive_chat_service.preview_outreach(payload.wecom_user_id) + channel, external_user_id = _resolve_channel_identity(payload.channel, payload.external_user_id, payload.wecom_user_id) + return await proactive_chat_service.preview_outreach(channel, external_user_id) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) @@ -84,18 +101,21 @@ async def preview_proactive_chat(payload: ProactiveChatActionRequest, _: bool = @router.post("/proactive-chat/run-once") async def run_proactive_chat(payload: ProactiveChatActionRequest, _: bool = Depends(require_admin)): try: - return await proactive_chat_service.run_outreach_once(payload.wecom_user_id) + channel, external_user_id = _resolve_channel_identity(payload.channel, payload.external_user_id, payload.wecom_user_id) + return await proactive_chat_service.run_outreach_once(channel, external_user_id) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.post("/persona/preview-prompt") async def preview_prompt(payload: PreviewRequest, _: bool = Depends(require_admin)): + channel, external_user_id = _resolve_channel_identity(payload.channel, payload.external_user_id, payload.wecom_user_id) preview = await run_preview_graph( { "preview_mode": "prompt", "user_message": payload.user_message, - "wecom_user_id": payload.wecom_user_id, + "channel": channel, + "external_user_id": external_user_id, "draft_config": payload.draft_config.model_dump() if payload.draft_config else None, } ) @@ -111,11 +131,13 @@ async def preview_prompt(payload: PreviewRequest, _: bool = Depends(require_admi @router.post("/persona/preview-reply") async def preview_reply(payload: PreviewRequest, _: bool = Depends(require_admin)): + channel, external_user_id = _resolve_channel_identity(payload.channel, payload.external_user_id, payload.wecom_user_id) preview = await run_preview_graph( { "preview_mode": "reply", "user_message": payload.user_message, - "wecom_user_id": payload.wecom_user_id, + "channel": channel, + "external_user_id": external_user_id, "draft_config": payload.draft_config.model_dump() if payload.draft_config else None, } ) @@ -132,34 +154,62 @@ async def preview_reply(payload: PreviewRequest, _: bool = Depends(require_admin @router.get("/users") async def list_users( - query: str = Query("", description="按企业微信用户 ID 或昵称搜索"), + query: str = Query("", description="按用户 ID 或昵称搜索"), limit: int = Query(20, ge=1, le=50), _: bool = Depends(require_admin), ): return {"items": await memory_service.list_users(query=query, limit=limit)} +@router.get("/users/{channel}/{external_user_id}/memory") +async def get_user_memory(channel: str, external_user_id: str, _: bool = Depends(require_admin)): + payload = await memory_service.get_user_memory(channel, external_user_id) + if not payload: + raise HTTPException(status_code=404, detail="User not found") + return payload + + +@router.put("/users/{channel}/{external_user_id}/memory") +async def update_user_memory( + channel: str, + external_user_id: str, + payload: UserMemoryPayload, + _: bool = Depends(require_admin), +): + return await memory_service.upsert_user_memory(channel, external_user_id, payload.model_dump()) + + +@router.get("/users/{channel}/{external_user_id}/conversations") +async def get_user_conversations( + channel: str, + external_user_id: str, + limit: int = Query(8, ge=1, le=20), + _: bool = Depends(require_admin), +): + return {"items": await memory_service.get_recent_conversations(channel, external_user_id, limit=limit)} + + @router.get("/users/{wecom_user_id}/memory") -async def get_user_memory(wecom_user_id: str, _: bool = Depends(require_admin)): - payload = await memory_service.get_user_memory(wecom_user_id) +async def get_user_memory_legacy(wecom_user_id: str, _: bool = Depends(require_admin)): + payload = await memory_service.get_user_memory("wecom", wecom_user_id) if not payload: raise HTTPException(status_code=404, detail="User not found") return payload @router.put("/users/{wecom_user_id}/memory") -async def update_user_memory( +async def update_user_memory_legacy( wecom_user_id: str, payload: UserMemoryPayload, _: bool = Depends(require_admin), ): - return await memory_service.upsert_user_memory(wecom_user_id, payload.model_dump()) + return await memory_service.upsert_user_memory("wecom", wecom_user_id, payload.model_dump()) @router.get("/users/{wecom_user_id}/conversations") -async def get_user_conversations( +async def get_user_conversations_legacy( wecom_user_id: str, limit: int = Query(8, ge=1, le=20), _: bool = Depends(require_admin), ): - return {"items": await memory_service.get_recent_conversations(wecom_user_id, limit=limit)} + return {"items": await memory_service.get_recent_conversations("wecom", wecom_user_id, limit=limit)} diff --git a/app/routers/setup.py b/app/routers/setup.py index 1d1372b..573fc92 100644 --- a/app/routers/setup.py +++ b/app/routers/setup.py @@ -4,7 +4,7 @@ from fastapi import APIRouter, HTTPException, Request -from app.schemas.admin import SetupAdminPayload, SetupModelPayload, SetupWeComPayload +from app.schemas.admin import SetupAdminPayload, SetupModelPayload, SetupNapCatPayload, SetupWeComPayload from app.services.runtime_config_service import runtime_config_service from app.services.setup_service import setup_service from app.services.tunnel_service import tunnel_service @@ -103,6 +103,19 @@ async def save_setup_wecom(payload: SetupWeComPayload, request: Request): return setup_service.get_status() +@router.put("/config/napcat") +async def save_setup_napcat(payload: SetupNapCatPayload, request: Request): + _require_setup_write_access(request) + runtime_config_service.save_section( + "napcat", + { + "ws_url": payload.ws_url.strip(), + "ws_token": payload.ws_token.strip(), + }, + ) + return setup_service.get_status() + + @router.put("/config/admin") async def save_setup_admin(payload: SetupAdminPayload, request: Request): _require_setup_write_access(request) diff --git a/app/routers/wecom.py b/app/routers/wecom.py index 178fa5f..1ce3198 100644 --- a/app/routers/wecom.py +++ b/app/routers/wecom.py @@ -87,7 +87,8 @@ async def wecom_callback_handler( await run_incoming_message_graph( { - "user_id": message.get("from_user"), + "channel": "wecom", + "external_user_id": message.get("from_user"), "user_content": message.get("content", ""), } ) diff --git a/app/schemas/admin.py b/app/schemas/admin.py index 9c9e202..bc352c1 100644 --- a/app/schemas/admin.py +++ b/app/schemas/admin.py @@ -35,6 +35,8 @@ class AgentPersonaPayload(BaseModel): class PreviewRequest(BaseModel): user_message: str + channel: Optional[str] = None + external_user_id: Optional[str] = None wecom_user_id: Optional[str] = None draft_config: Optional[AgentPersonaPayload] = None @@ -63,6 +65,8 @@ class QuietHoursPayload(BaseModel): class ProactiveChatPayload(BaseModel): enabled: bool = False + target_channel: str = "wecom" + target_external_user_id: str = "" target_wecom_user_id: str = "" scheduled_windows: List[ScheduledWindowPayload] = Field(default_factory=list) inactivity_trigger_hours: int = 6 @@ -73,6 +77,8 @@ class ProactiveChatPayload(BaseModel): class ProactiveChatActionRequest(BaseModel): + channel: Optional[str] = None + external_user_id: Optional[str] = None wecom_user_id: Optional[str] = None @@ -103,5 +109,10 @@ class SetupWeComPayload(BaseModel): public_base_url: str = "" +class SetupNapCatPayload(BaseModel): + ws_url: str = "" + ws_token: str = "" + + class SetupAdminPayload(BaseModel): password: str = "" diff --git a/app/services/__init__.py b/app/services/__init__.py index d6f96c5..b242321 100644 --- a/app/services/__init__.py +++ b/app/services/__init__.py @@ -17,6 +17,10 @@ "proactive_chat_service", "RuntimeConfigService", "runtime_config_service", + "NapCatService", + "napcat_service", + "ChannelDispatcher", + "channel_dispatcher", "TunnelService", "tunnel_service", ] @@ -35,6 +39,10 @@ "proactive_chat_service": ("app.services.proactive_chat_service", "proactive_chat_service"), "RuntimeConfigService": ("app.services.runtime_config_service", "RuntimeConfigService"), "runtime_config_service": ("app.services.runtime_config_service", "runtime_config_service"), + "NapCatService": ("app.services.napcat_service", "NapCatService"), + "napcat_service": ("app.services.napcat_service", "napcat_service"), + "ChannelDispatcher": ("app.services.channel_dispatcher", "ChannelDispatcher"), + "channel_dispatcher": ("app.services.channel_dispatcher", "channel_dispatcher"), "TunnelService": ("app.services.tunnel_service", "TunnelService"), "tunnel_service": ("app.services.tunnel_service", "tunnel_service"), } diff --git a/app/services/channel_dispatcher.py b/app/services/channel_dispatcher.py new file mode 100644 index 0000000..91cee46 --- /dev/null +++ b/app/services/channel_dispatcher.py @@ -0,0 +1,29 @@ +""" +Channel-aware outbound dispatcher. +""" + +from __future__ import annotations + +from typing import Dict + +from app.services.wecom_service import wecom_service + + +class ChannelDispatcher: + async def send_text(self, channel: str, external_user_id: str, content: str) -> Dict[str, object]: + lowered = (channel or "").strip().lower() + if lowered == "wecom": + await wecom_service.send_text_message(external_user_id, content) + return {"channel": "wecom", "status": "sent"} + + if lowered == "napcat": + from app.services.napcat_service import napcat_service + + await napcat_service.send_private_text(external_user_id, content) + return {"channel": "napcat", "status": "sent"} + + raise ValueError(f"Unsupported channel: {channel}") + + +channel_dispatcher = ChannelDispatcher() + diff --git a/app/services/llm_service.py b/app/services/llm_service.py index 7c4eec7..2954d70 100644 --- a/app/services/llm_service.py +++ b/app/services/llm_service.py @@ -254,7 +254,11 @@ async def maybe_collect_web_context(self, user_message: str) -> Dict[str, object if not query: return {"enabled": config["zhipu_web_search_enabled"], "triggered": False, "query": "", "results": []} - results = await self.web_search(query) + try: + results = await self.web_search(query) + except Exception as exc: + print(f"web search unavailable, fallback to no-search context: {exc}") + results = [] return { "enabled": config["zhipu_web_search_enabled"], "triggered": bool(results), diff --git a/app/services/memory_service.py b/app/services/memory_service.py index db05c22..b769a1b 100644 --- a/app/services/memory_service.py +++ b/app/services/memory_service.py @@ -31,12 +31,18 @@ def __init__(self): self._background_tasks: set[asyncio.Task] = set() self._user_locks: Dict[str, asyncio.Lock] = {} - async def get_or_create_user(self, wecom_user_id: str) -> Dict: + def _resolve_identity(self, channel: str, external_user_id: Optional[str] = None) -> tuple[str, str]: + if external_user_id is None: + return "wecom", channel + return channel, external_user_id + + async def get_or_create_user(self, channel: str, external_user_id: Optional[str] = None) -> Dict: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, wecom_user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: - user = self._create_user(db, wecom_user_id) + user = self._create_user(db, channel, external_user_id) user.last_interaction = datetime.now() user.total_conversations += 1 @@ -49,10 +55,11 @@ async def get_or_create_user(self, wecom_user_id: str) -> Dict: finally: db.close() - async def get_conversation_context(self, user_id: str) -> Dict: + async def get_conversation_context(self, channel: str, external_user_id: Optional[str] = None) -> Dict: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return {} @@ -90,10 +97,11 @@ async def get_conversation_context(self, user_id: str) -> Dict: finally: db.close() - async def get_recent_messages(self, user_id: str, limit: int = 10) -> List[Dict]: + async def get_recent_messages(self, channel: str, external_user_id: Optional[str] = None, limit: int = 10) -> List[Dict]: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return [] @@ -110,10 +118,11 @@ async def get_recent_messages(self, user_id: str, limit: int = 10) -> List[Dict] finally: db.close() - async def get_recent_agent_replies(self, user_id: str, limit: int = 3) -> List[str]: + async def get_recent_agent_replies(self, channel: str, external_user_id: Optional[str] = None, limit: int = 3) -> List[str]: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return [] @@ -125,16 +134,23 @@ async def get_recent_agent_replies(self, user_id: str, limit: int = 3) -> List[s async def save_conversation( self, - user_id: str, + *, + channel: Optional[str] = None, + external_user_id: Optional[str] = None, user_message: str, agent_message: str, user_emotion: Dict, agent_emotion: Dict, memories_used: Optional[Dict] = None, + user_id: Optional[str] = None, ) -> Optional[int]: + if user_id and not external_user_id: + channel, external_user_id = "wecom", user_id + if not channel or not external_user_id: + return None db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return None @@ -159,10 +175,17 @@ async def save_conversation( finally: db.close() - async def update_user_profile(self, user_id: str, profile_update: Dict) -> None: + async def update_user_profile( + self, + channel: str, + external_user_id: Optional[str] = None, + profile_update: Optional[Dict] = None, + ) -> None: + channel, external_user_id = self._resolve_identity(channel, external_user_id) + profile_update = profile_update or {} db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return @@ -183,7 +206,7 @@ async def list_users(self, query: str = "", limit: int = 20) -> List[Dict]: if cleaned_query: like_value = f"%{cleaned_query}%" user_query = user_query.filter( - (User.wecom_user_id.ilike(like_value)) | (User.nickname.ilike(like_value)) + (User.external_user_id.ilike(like_value)) | (User.nickname.ilike(like_value)) ) users = ( @@ -195,7 +218,8 @@ async def list_users(self, query: str = "", limit: int = 20) -> List[Dict]: return [ { - "wecom_user_id": user.wecom_user_id, + "channel": user.channel, + "external_user_id": user.external_user_id, "nickname": user.nickname or "", "avatar_url": user.avatar_url or "", "total_conversations": user.total_conversations, @@ -207,10 +231,11 @@ async def list_users(self, query: str = "", limit: int = 20) -> List[Dict]: finally: db.close() - async def get_user_memory(self, wecom_user_id: str, query_text: str = "") -> Optional[Dict]: + async def get_user_memory(self, channel: str, external_user_id: Optional[str] = None, query_text: str = "") -> Optional[Dict]: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, wecom_user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return None @@ -228,12 +253,14 @@ async def get_user_memory(self, wecom_user_id: str, query_text: str = "") -> Opt finally: db.close() - async def upsert_user_memory(self, wecom_user_id: str, payload: Dict) -> Dict: + async def upsert_user_memory(self, channel: str, external_user_id: Optional[str] = None, payload: Dict = None) -> Dict: + channel, external_user_id = self._resolve_identity(channel, external_user_id) + payload = payload or {} db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, wecom_user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: - user = self._create_user(db, wecom_user_id) + user = self._create_user(db, channel, external_user_id) user.nickname = str(payload.get("nickname") or "").strip() or None user.avatar_url = str(payload.get("avatar_url") or "").strip() or None @@ -251,13 +278,14 @@ async def upsert_user_memory(self, wecom_user_id: str, payload: Dict) -> Dict: finally: db.close() - memory = await self.get_user_memory(wecom_user_id) - return memory or self._empty_user_memory(wecom_user_id) + memory = await self.get_user_memory(channel, external_user_id) + return memory or self._empty_user_memory(channel, external_user_id) - async def get_recent_conversations(self, wecom_user_id: str, limit: int = 8) -> List[Dict]: + async def get_recent_conversations(self, channel: str, external_user_id: Optional[str] = None, limit: int = 8) -> List[Dict]: + channel, external_user_id = self._resolve_identity(channel, external_user_id) db = SessionLocal() try: - user = self._get_user_by_wecom_id(db, wecom_user_id) + user = self._get_user_by_channel_external_id(db, channel, external_user_id) if not user: return [] conversations = self._get_recent_conversations_rows(db, user.id, limit=limit) @@ -267,13 +295,20 @@ async def get_recent_conversations(self, wecom_user_id: str, limit: int = 8) -> def schedule_memory_processing( self, - wecom_user_id: str, + *, + channel: str = "wecom", + external_user_id: Optional[str] = None, conversation_id: Optional[int], user_message: str, agent_message: str, user_emotion: Dict, agent_emotion: Dict, + wecom_user_id: Optional[str] = None, ) -> None: + if wecom_user_id and not external_user_id: + channel, external_user_id = "wecom", wecom_user_id + else: + channel, external_user_id = self._resolve_identity(channel, external_user_id) if not conversation_id: return @@ -284,7 +319,8 @@ def schedule_memory_processing( task = loop.create_task( self.process_memory_update( - wecom_user_id=wecom_user_id, + channel=channel, + external_user_id=external_user_id, conversation_id=conversation_id, user_message=user_message, agent_message=agent_message, @@ -297,20 +333,29 @@ def schedule_memory_processing( async def process_memory_update( self, - wecom_user_id: str, + *, + channel: str = "wecom", + external_user_id: Optional[str] = None, conversation_id: int, user_message: str, agent_message: str, user_emotion: Dict, agent_emotion: Dict, + wecom_user_id: Optional[str] = None, ) -> None: - lock = self._user_locks.setdefault(wecom_user_id, asyncio.Lock()) + if wecom_user_id and not external_user_id: + channel, external_user_id = "wecom", wecom_user_id + else: + channel, external_user_id = self._resolve_identity(channel, external_user_id) + user_key = self.build_user_key(channel, external_user_id) + lock = self._user_locks.setdefault(user_key, asyncio.Lock()) async with lock: from app.graph import run_memory_update_graph await run_memory_update_graph( { - "wecom_user_id": wecom_user_id, + "channel": channel, + "external_user_id": external_user_id, "conversation_id": conversation_id, "user_message": user_message, "agent_message": agent_message, @@ -326,12 +371,17 @@ def _on_background_task_done(self, task: asyncio.Task) -> None: except Exception as exc: print(f"记忆异步提炼失败: {exc}") - def _get_user_by_wecom_id(self, db, wecom_user_id: str) -> Optional[User]: - return db.query(User).filter(User.wecom_user_id == wecom_user_id).first() + def _get_user_by_channel_external_id(self, db, channel: str, external_user_id: str) -> Optional[User]: + return ( + db.query(User) + .filter(User.channel == channel, User.external_user_id == external_user_id) + .first() + ) - def _create_user(self, db, wecom_user_id: str) -> User: + def _create_user(self, db, channel: str, external_user_id: str) -> User: user = User( - wecom_user_id=wecom_user_id, + channel=channel, + external_user_id=external_user_id, profile=self._get_default_profile(), basic_info={}, emotional_patterns={}, @@ -345,6 +395,9 @@ def _create_user(self, db, wecom_user_id: str) -> User: db.refresh(user) return user + def build_user_key(self, channel: str, external_user_id: str) -> str: + return f"{channel}:{external_user_id}" + def _get_or_create_short_term_memory(self, db, user_id: int) -> ShortTermMemory: memory = db.query(ShortTermMemory).filter(ShortTermMemory.user_id == user_id).first() if not memory: @@ -421,7 +474,8 @@ def _build_profile_snapshot( ) -> Dict: profile = user.profile or self._get_default_profile() return { - "wecom_user_id": user.wecom_user_id, + "channel": user.channel, + "external_user_id": user.external_user_id, "nickname": user.nickname or "", "avatar_url": user.avatar_url or "", "profile": profile, @@ -462,9 +516,10 @@ def _serialize_memory_item(self, item: MemoryItem) -> Dict: "source_conversation_id": item.source_conversation_id, } - def _empty_user_memory(self, wecom_user_id: str) -> Dict: + def _empty_user_memory(self, channel: str, external_user_id: str) -> Dict: return { - "wecom_user_id": wecom_user_id, + "channel": channel, + "external_user_id": external_user_id, "nickname": "", "avatar_url": "", "basic_info": {}, diff --git a/app/services/napcat_service.py b/app/services/napcat_service.py new file mode 100644 index 0000000..99b5018 --- /dev/null +++ b/app/services/napcat_service.py @@ -0,0 +1,132 @@ +""" +NapCat OneBot11 forward WebSocket client service. +""" + +from __future__ import annotations + +import asyncio +import json +import random +from typing import Optional + +from app.config import settings +from app.services.runtime_config_service import runtime_config_service + +try: + import websockets +except Exception: # pragma: no cover + websockets = None # type: ignore[assignment] + + +class NapCatService: + def __init__(self) -> None: + self._task: Optional[asyncio.Task] = None + self._ws = None + self._send_lock = asyncio.Lock() + self._stop_event = asyncio.Event() + + def _config(self) -> dict: + runtime = runtime_config_service.get_effective_napcat_config() + return { + "ws_url": runtime.get("ws_url") or settings.napcat_ws_url, + "ws_token": runtime.get("ws_token") or settings.napcat_ws_token, + } + + async def start(self) -> None: + if websockets is None: + print("NapCat disabled: websockets is unavailable") + return + if self._task and not self._task.done(): + return + cfg = self._config() + if not str(cfg["ws_url"]).strip(): + print("NapCat disabled: NAPCAT_WS_URL is empty") + return + self._stop_event.clear() + self._task = asyncio.create_task(self._run_loop()) + + async def stop(self) -> None: + self._stop_event.set() + task = self._task + self._task = None + if task: + task.cancel() + try: + await task + except asyncio.CancelledError: + pass + if self._ws: + try: + await self._ws.close() + except Exception: + pass + self._ws = None + + async def _run_loop(self) -> None: + delay = 1.0 + while not self._stop_event.is_set(): + cfg = self._config() + ws_url = str(cfg["ws_url"]).strip() + if not ws_url: + return + headers = {} + token = str(cfg.get("ws_token") or "").strip() + if token: + headers["Authorization"] = f"Bearer {token}" + try: + async with websockets.connect(ws_url, additional_headers=headers) as ws: + self._ws = ws + delay = 1.0 + async for message in ws: + await self._handle_message(message) + except asyncio.CancelledError: + raise + except Exception as exc: + print(f"NapCat connection error: {exc}") + finally: + self._ws = None + + sleep_seconds = min(30.0, delay + random.uniform(0, 0.8)) + await asyncio.sleep(sleep_seconds) + delay = min(30.0, delay * 2) + + async def _handle_message(self, message: str) -> None: + try: + payload = json.loads(message) + except json.JSONDecodeError: + return + + if payload.get("post_type") != "message": + return + if payload.get("message_type") != "private": + return + + content = str(payload.get("raw_message") or "").strip() + external_user_id = str(payload.get("user_id") or "").strip() + if not content or not external_user_id: + return + + from app.graph import run_incoming_message_graph + + await run_incoming_message_graph( + { + "channel": "napcat", + "external_user_id": external_user_id, + "user_content": content, + } + ) + + async def send_private_text(self, external_user_id: str, content: str) -> None: + if not self._ws: + raise RuntimeError("NapCat websocket is not connected") + payload = { + "action": "send_private_msg", + "params": {"user_id": int(external_user_id) if external_user_id.isdigit() else external_user_id, "message": content}, + "echo": f"lovagent-{int(asyncio.get_running_loop().time() * 1000)}", + } + async with self._send_lock: + await self._ws.send(json.dumps(payload, ensure_ascii=False)) + + +napcat_service = NapCatService() + diff --git a/app/services/proactive_chat_service.py b/app/services/proactive_chat_service.py index 3c1c2a3..b3882a9 100644 --- a/app/services/proactive_chat_service.py +++ b/app/services/proactive_chat_service.py @@ -1,27 +1,17 @@ -""" -主动聊天服务 +""" +Proactive chat service. """ import asyncio from datetime import datetime, timedelta -from typing import Dict, List, Optional +from typing import Dict, List, Optional, Tuple from sqlalchemy.exc import OperationalError from app.config import settings from app.models.admin import ProactiveChatConfig, ProactiveChatLog from app.models.database import SessionLocal -from app.models.user import Conversation, User -from app.prompts.templates import ( - build_morning_greeting, - build_night_greeting, - build_proactive_prompt, -) -from app.services.llm_service import glm_service -from app.services.memory_service import memory_service -from app.services.persona_service import persona_service -from app.services.wecom_service import wecom_service -from app.utils.helpers import get_response_constraints, get_time_period +from app.models.user import User DEFAULT_PROACTIVE_CHAT_CONFIG_KEY = "default_proactive_chat" @@ -33,7 +23,8 @@ DEFAULT_QUIET_HOURS = {"enabled": True, "start": "23:00", "end": "09:00"} DEFAULT_PROACTIVE_CHAT_CONFIG = { "enabled": False, - "target_wecom_user_id": "", + "target_channel": "wecom", + "target_external_user_id": "", "scheduled_windows": DEFAULT_SCHEDULED_WINDOWS, "inactivity_trigger_hours": 6, "quiet_hours": DEFAULT_QUIET_HOURS, @@ -44,7 +35,7 @@ class ProactiveChatService: - """主动聊天配置、调度和发送服务。""" + """Proactive chat config, scheduling and dispatch.""" def __init__(self): self.scheduler_interval_seconds = max(30, settings.proactive_scheduler_interval_seconds) @@ -67,7 +58,8 @@ def get_config(self) -> Dict: return self._build_payload( { "enabled": record.enabled, - "target_wecom_user_id": record.target_wecom_user_id or "", + "target_channel": record.target_channel or "wecom", + "target_external_user_id": record.target_external_user_id or "", "scheduled_windows": record.scheduled_windows or [], "inactivity_trigger_hours": record.inactivity_trigger_hours, "quiet_hours": record.quiet_hours or {}, @@ -99,7 +91,8 @@ def save_config(self, config: Dict) -> Dict: db.add(record) record.enabled = normalized["enabled"] - record.target_wecom_user_id = normalized["target_wecom_user_id"] or None + record.target_channel = normalized["target_channel"] + record.target_external_user_id = normalized["target_external_user_id"] or None record.scheduled_windows = normalized["scheduled_windows"] record.inactivity_trigger_hours = normalized["inactivity_trigger_hours"] record.quiet_hours = normalized["quiet_hours"] @@ -113,13 +106,14 @@ def save_config(self, config: Dict) -> Dict: finally: db.close() - async def preview_outreach(self, wecom_user_id: Optional[str] = None) -> Dict: - target_wecom_user_id = self._resolve_target_user_id(wecom_user_id) + async def preview_outreach(self, channel: Optional[str] = None, external_user_id: Optional[str] = None) -> Dict: + target_channel, target_external_user_id = self._resolve_target_user(channel, external_user_id) from app.graph import run_proactive_chat_graph payload = await run_proactive_chat_graph( { - "target_wecom_user_id": target_wecom_user_id, + "target_channel": target_channel, + "target_external_user_id": target_external_user_id, "trigger_type": "manual", "window_key": None, "send_delivery": False, @@ -127,13 +121,14 @@ async def preview_outreach(self, wecom_user_id: Optional[str] = None) -> Dict: ) return self._format_graph_payload(payload) - async def run_outreach_once(self, wecom_user_id: Optional[str] = None) -> Dict: - target_wecom_user_id = self._resolve_target_user_id(wecom_user_id) + async def run_outreach_once(self, channel: Optional[str] = None, external_user_id: Optional[str] = None) -> Dict: + target_channel, target_external_user_id = self._resolve_target_user(channel, external_user_id) from app.graph import run_proactive_chat_graph payload = await run_proactive_chat_graph( { - "target_wecom_user_id": target_wecom_user_id, + "target_channel": target_channel, + "target_external_user_id": target_external_user_id, "trigger_type": "manual", "window_key": None, "send_delivery": True, @@ -151,7 +146,8 @@ async def dispatch_due_messages(self) -> Optional[Dict]: payload = await run_proactive_chat_graph( { - "target_wecom_user_id": config["target_wecom_user_id"], + "target_channel": config["target_channel"], + "target_external_user_id": config["target_external_user_id"], "trigger_type": due["trigger_type"], "window_key": due.get("window_key"), "send_delivery": True, @@ -164,7 +160,7 @@ async def scheduler_loop(self) -> None: try: result = await self.dispatch_due_messages() if result and result.get("delivery", {}).get("status") == "sent": - print(f"主动聊天发送成功: {result['target_wecom_user_id']}") + print(f"主动聊天发送成功: {result['target_channel']}:{result['target_external_user_id']}") except asyncio.CancelledError: raise except Exception as exc: @@ -175,11 +171,16 @@ async def scheduler_loop(self) -> None: def _build_payload(self, config: Dict, updated_at: Optional[str] = None) -> Dict: normalized = self._normalize_config(config) normalized["updated_at"] = updated_at + normalized["target_wecom_user_id"] = ( + normalized["target_external_user_id"] if normalized["target_channel"] == "wecom" else "" + ) return normalized def _format_graph_payload(self, payload: Dict) -> Dict: return { - "target_wecom_user_id": payload["target_wecom_user_id"], + "target_channel": payload["target_channel"], + "target_external_user_id": payload["target_external_user_id"], + "target_wecom_user_id": payload["target_external_user_id"] if payload["target_channel"] == "wecom" else "", "trigger_type": payload["trigger_type"], "window_key": payload.get("window_key"), "prompt": payload.get("prompt", ""), @@ -199,7 +200,9 @@ def _normalize_config(self, config: Optional[Dict]) -> Dict: return merged merged["enabled"] = bool(config.get("enabled", merged["enabled"])) - merged["target_wecom_user_id"] = str(config.get("target_wecom_user_id") or "").strip() + legacy_target = str(config.get("target_wecom_user_id") or "").strip() + merged["target_channel"] = str(config.get("target_channel") or "wecom").strip() or "wecom" + merged["target_external_user_id"] = str(config.get("target_external_user_id") or legacy_target).strip() merged["scheduled_windows"] = self._normalize_scheduled_windows(config.get("scheduled_windows")) merged["quiet_hours"] = self._normalize_quiet_hours(config.get("quiet_hours")) merged["tone_hint"] = str(config.get("tone_hint") or merged["tone_hint"]).strip() or merged["tone_hint"] @@ -275,176 +278,34 @@ def _coerce_int(self, value: object, fallback: int, minimum: int, maximum: int) number = fallback return max(minimum, min(maximum, number)) - def _resolve_target_user_id(self, wecom_user_id: Optional[str]) -> str: - target_wecom_user_id = str(wecom_user_id or "").strip() - if target_wecom_user_id: - return target_wecom_user_id + def _resolve_target_user(self, channel: Optional[str], external_user_id: Optional[str]) -> Tuple[str, str]: + target_channel = str(channel or "").strip() + target_external_user_id = str(external_user_id or "").strip() + if target_channel and target_external_user_id: + return target_channel, target_external_user_id config = self.get_config() - target_wecom_user_id = config["target_wecom_user_id"] - if not target_wecom_user_id: + target_channel = config["target_channel"] + target_external_user_id = config["target_external_user_id"] + if not target_external_user_id: raise ValueError("请先配置主动聊天目标用户") - return target_wecom_user_id - - async def _build_outreach_payload( - self, - target_wecom_user_id: str, - trigger_type: str, - window_key: Optional[str] = None, - ) -> Dict: - persona_config = persona_service.get_persona_config() - proactive_config = self.get_config() - user_memory = await memory_service.get_user_memory(target_wecom_user_id) - context = await memory_service.get_conversation_context(target_wecom_user_id) - recent_agent_replies = await memory_service.get_recent_agent_replies(target_wecom_user_id, limit=3) - prompt = build_proactive_prompt( - trigger_type=trigger_type, - current_time=datetime.now().strftime("%Y-%m-%d %H:%M:%S"), - persona_config=persona_config, - user_profile=user_memory, - context=context, - recent_agent_replies=recent_agent_replies, - tone_hint=proactive_config["tone_hint"], - ) - - response_constraints = get_response_constraints( - "我想主动开启一段自然微信聊天", - persona_config.get("response_preferences"), - ) - - reply = "" - try: - reply = await glm_service.chat_with_context( - system_prompt=prompt, - user_message="请直接输出一条现在要主动发给他的微信消息。", - context_messages=[], - temperature=0.92, - top_p=0.95, - max_tokens=int(response_constraints["max_tokens"]), - task_type="proactive", - ) - except Exception as exc: - print(f"生成主动聊天文案失败,使用兜底: {exc}") - - if not reply: - reply = self._build_fallback_message(target_wecom_user_id, trigger_type, user_memory) - - return { - "target_wecom_user_id": target_wecom_user_id, - "trigger_type": trigger_type, - "window_key": window_key, - "prompt": prompt, - "reply": reply, - "persona_config": persona_config, - "user_memory": user_memory, - "config": proactive_config, - } + return target_channel, target_external_user_id def _build_fallback_message( self, - target_wecom_user_id: str, + target_channel: str, + target_external_user_id: str, trigger_type: str, user_memory: Optional[Dict], ) -> str: - nickname = str((user_memory or {}).get("nickname") or "").strip() - prefix = f"{nickname}," if nickname else "" - time_period = get_time_period() - - if trigger_type == "scheduled" and time_period == "早晨": - return build_morning_greeting() - if trigger_type == "scheduled" and time_period in {"夜晚", "深夜"}: - return build_night_greeting() + _ = (target_channel, target_external_user_id, user_memory) if trigger_type == "inactivity": - return f"{prefix}刚刚突然想到你了,今天过得怎么样呀?" - return f"{prefix}刚刚想起你,想来和你说句话,在忙吗?" - - async def _deliver_outreach( - self, - target_wecom_user_id: str, - trigger_type: str, - window_key: Optional[str], - content: str, - ) -> Dict: - sent_at = datetime.now() - status = "failed" - error_message = None - - try: - await wecom_service.send_text_message(target_wecom_user_id, content) - status = "sent" - self._save_proactive_conversation(target_wecom_user_id, content, sent_at) - except Exception as exc: - error_message = str(exc) - print(f"主动聊天发送失败: {exc}") - finally: - self._save_log( - target_wecom_user_id=target_wecom_user_id, - trigger_type=trigger_type, - window_key=window_key, - content=content, - status=status, - error_message=error_message, - sent_at=sent_at, - ) - - return { - "attempted": True, - "status": status, - "error_message": error_message, - "sent_at": sent_at.isoformat(), - } - - def _save_proactive_conversation(self, target_wecom_user_id: str, content: str, sent_at: datetime) -> None: - db = SessionLocal() - try: - user = db.query(User).filter(User.wecom_user_id == target_wecom_user_id).first() - if not user: - return - - conversation = Conversation( - user_id=user.id, - user_message="", - agent_message=content, - user_emotion=None, - agent_emotion="happy", - agent_emotion_intensity=40, - context_used=True, - memories_used={"source": "proactive"}, - created_at=sent_at, - ) - db.add(conversation) - db.commit() - finally: - db.close() - - def _save_log( - self, - target_wecom_user_id: str, - trigger_type: str, - window_key: Optional[str], - content: str, - status: str, - error_message: Optional[str], - sent_at: datetime, - ) -> None: - db = SessionLocal() - try: - log = ProactiveChatLog( - target_wecom_user_id=target_wecom_user_id, - trigger_type=trigger_type, - window_key=window_key, - content=content, - status=status, - error_message=error_message, - sent_at=sent_at, - ) - db.add(log) - db.commit() - finally: - db.close() + return "刚刚突然想到你了,今天过得怎么样呀?" + return "刚刚想起你,想来和你说句话,在忙吗?" def _resolve_due_trigger(self, config: Dict) -> Optional[Dict]: - if not config.get("enabled") or not config.get("target_wecom_user_id"): + config = self._normalize_config(config) + if not config.get("enabled") or not config.get("target_external_user_id"): return None now = datetime.now() @@ -453,7 +314,14 @@ def _resolve_due_trigger(self, config: Dict) -> Optional[Dict]: db = SessionLocal() try: - user = db.query(User).filter(User.wecom_user_id == config["target_wecom_user_id"]).first() + user = ( + db.query(User) + .filter( + User.channel == config["target_channel"], + User.external_user_id == config["target_external_user_id"], + ) + .first() + ) if not user: return None @@ -463,10 +331,10 @@ def _resolve_due_trigger(self, config: Dict) -> Optional[Dict]: if minutes_since_last_interaction < config["min_interval_minutes"]: return None - if self._count_sent_today(db, config["target_wecom_user_id"], now) >= config["max_messages_per_day"]: + if self._count_sent_today(db, config["target_channel"], config["target_external_user_id"], now) >= config["max_messages_per_day"]: return None - last_sent = self._get_last_success_log(db, config["target_wecom_user_id"]) + last_sent = self._get_last_success_log(db, config["target_channel"], config["target_external_user_id"]) if last_sent and (now - last_sent.sent_at).total_seconds() / 60 < config["min_interval_minutes"]: return None @@ -475,7 +343,8 @@ def _resolve_due_trigger(self, config: Dict) -> Optional[Dict]: continue if self._is_window_due(now, window["time"]) and not self._window_sent_today( db, - config["target_wecom_user_id"], + config["target_channel"], + config["target_external_user_id"], window["key"], now, ): @@ -484,7 +353,12 @@ def _resolve_due_trigger(self, config: Dict) -> Optional[Dict]: if ( last_interaction and now - last_interaction >= timedelta(hours=config["inactivity_trigger_hours"]) - and not self._already_sent_inactivity_since(db, config["target_wecom_user_id"], last_interaction) + and not self._already_sent_inactivity_since( + db, + config["target_channel"], + config["target_external_user_id"], + last_interaction, + ) ): return {"trigger_type": "inactivity", "window_key": None} finally: @@ -510,13 +384,14 @@ def _is_in_quiet_hours(self, now: datetime, quiet_hours: Dict) -> bool: return start <= current < end return current >= start or current < end - def _count_sent_today(self, db, target_wecom_user_id: str, now: datetime) -> int: + def _count_sent_today(self, db, target_channel: str, target_external_user_id: str, now: datetime) -> int: start_of_day = now.replace(hour=0, minute=0, second=0, microsecond=0) end_of_day = start_of_day + timedelta(days=1) return ( db.query(ProactiveChatLog) .filter( - ProactiveChatLog.target_wecom_user_id == target_wecom_user_id, + ProactiveChatLog.target_channel == target_channel, + ProactiveChatLog.target_external_user_id == target_external_user_id, ProactiveChatLog.status == "sent", ProactiveChatLog.sent_at >= start_of_day, ProactiveChatLog.sent_at < end_of_day, @@ -524,13 +399,14 @@ def _count_sent_today(self, db, target_wecom_user_id: str, now: datetime) -> int .count() ) - def _window_sent_today(self, db, target_wecom_user_id: str, window_key: str, now: datetime) -> bool: + def _window_sent_today(self, db, target_channel: str, target_external_user_id: str, window_key: str, now: datetime) -> bool: start_of_day = now.replace(hour=0, minute=0, second=0, microsecond=0) end_of_day = start_of_day + timedelta(days=1) return ( db.query(ProactiveChatLog) .filter( - ProactiveChatLog.target_wecom_user_id == target_wecom_user_id, + ProactiveChatLog.target_channel == target_channel, + ProactiveChatLog.target_external_user_id == target_external_user_id, ProactiveChatLog.status == "sent", ProactiveChatLog.trigger_type == "scheduled", ProactiveChatLog.window_key == window_key, @@ -541,11 +417,12 @@ def _window_sent_today(self, db, target_wecom_user_id: str, window_key: str, now is not None ) - def _already_sent_inactivity_since(self, db, target_wecom_user_id: str, since: datetime) -> bool: + def _already_sent_inactivity_since(self, db, target_channel: str, target_external_user_id: str, since: datetime) -> bool: return ( db.query(ProactiveChatLog) .filter( - ProactiveChatLog.target_wecom_user_id == target_wecom_user_id, + ProactiveChatLog.target_channel == target_channel, + ProactiveChatLog.target_external_user_id == target_external_user_id, ProactiveChatLog.status == "sent", ProactiveChatLog.trigger_type == "inactivity", ProactiveChatLog.sent_at >= since, @@ -554,11 +431,12 @@ def _already_sent_inactivity_since(self, db, target_wecom_user_id: str, since: d is not None ) - def _get_last_success_log(self, db, target_wecom_user_id: str) -> Optional[ProactiveChatLog]: + def _get_last_success_log(self, db, target_channel: str, target_external_user_id: str) -> Optional[ProactiveChatLog]: return ( db.query(ProactiveChatLog) .filter( - ProactiveChatLog.target_wecom_user_id == target_wecom_user_id, + ProactiveChatLog.target_channel == target_channel, + ProactiveChatLog.target_external_user_id == target_external_user_id, ProactiveChatLog.status == "sent", ) .order_by(ProactiveChatLog.sent_at.desc(), ProactiveChatLog.id.desc()) diff --git a/app/services/runtime_config_service.py b/app/services/runtime_config_service.py index 323a9ff..23b6b35 100644 --- a/app/services/runtime_config_service.py +++ b/app/services/runtime_config_service.py @@ -38,6 +38,10 @@ "token": "", "encoding_aes_key": "", }, + "napcat": { + "ws_url": "", + "ws_token": "", + }, "deployment": { "public_base_url": "", }, @@ -189,6 +193,13 @@ def get_effective_wecom_config(self) -> Dict: "encoding_aes_key": config["encoding_aes_key"] or settings.wecom_encoding_aes_key, } + def get_effective_napcat_config(self) -> Dict: + config = self.get_config()["napcat"] + return { + "ws_url": str(config.get("ws_url") or settings.napcat_ws_url).strip(), + "ws_token": str(config.get("ws_token") or settings.napcat_ws_token).strip(), + } + def get_effective_public_base_url(self) -> str: deployment = self.get_config()["deployment"] return str(deployment.get("public_base_url") or settings.public_base_url).strip() @@ -224,6 +235,7 @@ def get_status_payload(self) -> Dict: raw = self.get_config() effective_model = self.get_effective_model_config() effective_wecom = self.get_effective_wecom_config() + effective_napcat = self.get_effective_napcat_config() effective_public_base_url = self.get_effective_public_base_url() effective_admin_password = self.get_effective_admin_password() @@ -240,6 +252,7 @@ def get_status_payload(self) -> Dict: effective_wecom["encoding_aes_key"], ] ), + "napcat_configured": bool(effective_napcat["ws_url"]), "admin_configured": bool(effective_admin_password), "deployment_configured": bool(effective_public_base_url), }, @@ -259,6 +272,8 @@ def get_status_payload(self) -> Dict: "has_wecom_secret": bool(effective_wecom["secret"]), "has_wecom_token": bool(effective_wecom["token"]), "has_wecom_encoding_aes_key": bool(effective_wecom["encoding_aes_key"]), + "napcat_ws_url": effective_napcat["ws_url"], + "has_napcat_ws_token": bool(effective_napcat["ws_token"]), "has_admin_password": bool(effective_admin_password), }, "raw": { @@ -279,6 +294,10 @@ def get_status_payload(self) -> Dict: "has_token": bool(raw["wecom"]["token"]), "has_encoding_aes_key": bool(raw["wecom"]["encoding_aes_key"]), }, + "napcat": { + "ws_url": raw["napcat"]["ws_url"], + "has_ws_token": bool(raw["napcat"]["ws_token"]), + }, "deployment": { "public_base_url": raw["deployment"]["public_base_url"], }, diff --git a/requirements.txt b/requirements.txt index a60ad79..31b707e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -12,3 +12,4 @@ pydantic-settings>=2.0.0 itsdangerous>=2.2.0 langchain-core>=0.3.0,<0.4.0 langgraph>=0.2.0,<0.4.0 +websockets>=12.0 From 4d1cd706dc054301b5bf54c993243f7f5b0206e0 Mon Sep 17 00:00:00 2001 From: Rizx <157099086+rizx-r@users.noreply.github.com> Date: Wed, 8 Apr 2026 02:25:50 +0800 Subject: [PATCH 2/5] feat(admin-ui): support channel-aware user and proactive flows --- admin-ui/src/App.tsx | 73 ++++++++------------- admin-ui/src/api.ts | 28 +++++--- admin-ui/src/components/MemoryDesk.tsx | 55 ++++------------ admin-ui/src/components/ProactiveStudio.tsx | 35 ++++++---- admin-ui/src/components/UserSidebar.tsx | 20 +++--- admin-ui/src/lib/proactiveChat.ts | 3 +- admin-ui/src/lib/userMemory.ts | 3 +- admin-ui/src/types.ts | 12 ++-- 8 files changed, 102 insertions(+), 127 deletions(-) diff --git a/admin-ui/src/App.tsx b/admin-ui/src/App.tsx index 4bba5e0..a16e1c6 100644 --- a/admin-ui/src/App.tsx +++ b/admin-ui/src/App.tsx @@ -1,4 +1,4 @@ -import { FormEvent, startTransition, useDeferredValue, useEffect, useState } from "react"; +import { FormEvent, startTransition, useDeferredValue, useEffect, useState } from "react"; import { api } from "./api"; import { LoginShell } from "./components/LoginShell"; @@ -46,6 +46,7 @@ function App() { const deferredUserQuery = useDeferredValue(userQuery); const [users, setUsers] = useState([]); const [usersLoading, setUsersLoading] = useState(false); + const [selectedChannel, setSelectedChannel] = useState("wecom"); const [selectedUserId, setSelectedUserId] = useState(""); const [memoryDraft, setMemoryDraft] = useState(EMPTY_MEMORY); const [memoryLoading, setMemoryLoading] = useState(false); @@ -99,7 +100,8 @@ function App() { setProactiveConfig(proactive); setUsers(userPayload.items); if (userPayload.items[0] && !selectedUserId) { - setSelectedUserId(userPayload.items[0].wecom_user_id); + setSelectedChannel(userPayload.items[0].channel); + setSelectedUserId(userPayload.items[0].external_user_id); } }) .catch((error: Error) => { @@ -123,9 +125,10 @@ function App() { return; } - const stillExists = payload.items.some((item) => item.wecom_user_id === selectedUserId); + const stillExists = payload.items.some((item) => item.channel === selectedChannel && item.external_user_id === selectedUserId); if (!stillExists) { - setSelectedUserId(payload.items[0].wecom_user_id); + setSelectedChannel(payload.items[0].channel); + setSelectedUserId(payload.items[0].external_user_id); } }) .catch((error: Error) => { @@ -134,7 +137,7 @@ function App() { .finally(() => { setUsersLoading(false); }); - }, [authenticated, deferredUserQuery, selectedUserId, setupStatus?.setup_completed]); + }, [authenticated, deferredUserQuery, selectedChannel, selectedUserId, setupStatus?.setup_completed]); useEffect(() => { if (!authenticated || !setupStatus?.setup_completed || !selectedUserId) { @@ -143,7 +146,7 @@ function App() { setMemoryLoading(true); void api - .getUserMemory(selectedUserId) + .getUserMemory(selectedChannel, selectedUserId) .then((payload) => { setMemoryDraft(normalizeUserMemory(payload)); }) @@ -153,7 +156,7 @@ function App() { .finally(() => { setMemoryLoading(false); }); - }, [authenticated, selectedUserId, setupStatus?.setup_completed]); + }, [authenticated, selectedChannel, selectedUserId, setupStatus?.setup_completed]); async function handleLogin(event: FormEvent) { event.preventDefault(); @@ -204,7 +207,8 @@ function App() { try { const payload = { user_message: previewMessage, - wecom_user_id: selectedUserId || undefined, + channel: selectedChannel || undefined, + external_user_id: selectedUserId || undefined, draft_config: personaConfig, }; const response = mode === "prompt" ? await api.previewPrompt(payload) : await api.previewReply(payload); @@ -228,7 +232,7 @@ function App() { setMemorySaving(true); try { - const saved = await api.saveUserMemory(selectedUserId, memoryDraft); + const saved = await api.saveUserMemory(selectedChannel, selectedUserId, memoryDraft); setMemoryDraft(normalizeUserMemory(saved)); setStatusMessage("用户记忆已保存。"); } catch (error) { @@ -259,8 +263,8 @@ function App() { const response = mode === "preview" - ? await api.previewProactiveChat(saved.target_wecom_user_id) - : await api.runProactiveChatOnce(saved.target_wecom_user_id); + ? await api.previewProactiveChat(saved.target_channel, saved.target_external_user_id) + : await api.runProactiveChatOnce(saved.target_channel, saved.target_external_user_id); setProactivePrompt(response.prompt); setProactiveReply(response.reply); @@ -366,7 +370,7 @@ function App() { } function updateProactiveField( - field: "enabled" | "target_wecom_user_id" | "tone_hint", + field: "enabled" | "target_channel" | "target_external_user_id" | "tone_hint", value: boolean | string, ) { setProactiveConfig((current) => @@ -424,7 +428,7 @@ function App() { function handleEnterAdmin() { window.history.replaceState({}, "", ADMIN_PATH); if (!authenticated) { - setStatusMessage("请先登录管理员后台后再继续编辑配置。"); + setStatusMessage("请先登录管理员后台。"); return; } setStatusMessage("环境校验完成,已进入管理后台。"); @@ -444,45 +448,20 @@ function App() { } if (!setupStatus.setup_completed) { - return ( - - ); + return ; } if (window.location.pathname === SETUP_PATH) { - return ( - - ); + return ; } if (!authenticated) { - return ( - - ); + return ; } return (
- void handleLogout()} - /> + void handleLogout()} /> { @@ -503,7 +482,10 @@ function App() { setUserQuery(value); }); }} - onSelectUser={setSelectedUserId} + onSelectUser={(channel, externalUserId) => { + setSelectedChannel(channel); + setSelectedUserId(externalUserId); + }} />
@@ -563,7 +545,10 @@ function App() { previewReply={proactiveReply} deliveryStatus={proactiveDeliveryStatus} onToggleEnabled={(value) => updateProactiveField("enabled", value)} - onTargetUserChange={(value) => updateProactiveField("target_wecom_user_id", value)} + onTargetUserChange={(channel, externalUserId) => { + updateProactiveField("target_channel", channel); + updateProactiveField("target_external_user_id", externalUserId); + }} onWindowToggle={(key, enabled) => updateProactiveWindow(key, { enabled })} onWindowTimeChange={(key, value) => updateProactiveWindow(key, { time: value })} onQuietHoursToggle={(value) => updateQuietHours("enabled", value)} diff --git a/admin-ui/src/api.ts b/admin-ui/src/api.ts index 304a703..f2a0503 100644 --- a/admin-ui/src/api.ts +++ b/admin-ui/src/api.ts @@ -57,7 +57,8 @@ export const api = { }, previewPrompt(payload: { user_message: string; - wecom_user_id?: string | null; + channel?: string | null; + external_user_id?: string | null; draft_config?: PersonaConfig; }): Promise { return request("/admin-api/persona/preview-prompt", { @@ -70,7 +71,8 @@ export const api = { }, previewReply(payload: { user_message: string; - wecom_user_id?: string | null; + channel?: string | null; + external_user_id?: string | null; draft_config?: PersonaConfig; }): Promise { return request("/admin-api/persona/preview-reply", { @@ -89,11 +91,11 @@ export const api = { params.set("limit", "30"); return request(`/admin-api/users?${params.toString()}`); }, - getUserMemory(wecomUserId: string): Promise { - return request(`/admin-api/users/${encodeURIComponent(wecomUserId)}/memory`); + getUserMemory(channel: string, externalUserId: string): Promise { + return request(`/admin-api/users/${encodeURIComponent(channel)}/${encodeURIComponent(externalUserId)}/memory`); }, - saveUserMemory(wecomUserId: string, payload: UserMemory): Promise { - return request(`/admin-api/users/${encodeURIComponent(wecomUserId)}/memory`, { + saveUserMemory(channel: string, externalUserId: string, payload: UserMemory): Promise { + return request(`/admin-api/users/${encodeURIComponent(channel)}/${encodeURIComponent(externalUserId)}/memory`, { method: "PUT", body: JSON.stringify(payload), }); @@ -107,19 +109,25 @@ export const api = { body: JSON.stringify(payload), }).then(normalizeProactiveChatConfig); }, - previewProactiveChat(wecomUserId?: string): Promise { + previewProactiveChat(channel?: string, externalUserId?: string): Promise { return request("/admin-api/proactive-chat/preview", { method: "POST", - body: JSON.stringify({ wecom_user_id: wecomUserId || undefined }), + body: JSON.stringify({ + channel: channel || undefined, + external_user_id: externalUserId || undefined, + }), }).then((response) => ({ ...response, config: normalizeProactiveChatConfig(response.config), })); }, - runProactiveChatOnce(wecomUserId?: string): Promise { + runProactiveChatOnce(channel?: string, externalUserId?: string): Promise { return request("/admin-api/proactive-chat/run-once", { method: "POST", - body: JSON.stringify({ wecom_user_id: wecomUserId || undefined }), + body: JSON.stringify({ + channel: channel || undefined, + external_user_id: externalUserId || undefined, + }), }).then((response) => ({ ...response, config: normalizeProactiveChatConfig(response.config), diff --git a/admin-ui/src/components/MemoryDesk.tsx b/admin-ui/src/components/MemoryDesk.tsx index 59665fe..6088d58 100644 --- a/admin-ui/src/components/MemoryDesk.tsx +++ b/admin-ui/src/components/MemoryDesk.tsx @@ -1,4 +1,4 @@ -import type { UserMemory } from "../types"; +import type { UserMemory } from "../types"; import { KeyValueEditor, TextListEditor } from "./Editors"; type MemoryDeskProps = { @@ -17,15 +17,17 @@ type MemoryDeskProps = { export function MemoryDesk(props: MemoryDeskProps) { const { draft, loading, saving, onMemoryFieldChange, onKeyValueChange, onMilestonesChange, onSave } = props; + const userLabel = draft.external_user_id ? `${draft.channel}:${draft.external_user_id}` : "未选择用户"; + return (

Selected User

-

{draft.wecom_user_id || "未选择用户"}

+

{userLabel}

-
@@ -35,52 +37,21 @@ export function MemoryDesk(props: MemoryDeskProps) {
- onKeyValueChange("basic_info", nextValue)} - /> - onKeyValueChange("emotional_patterns", nextValue)} - /> - onKeyValueChange("preferences", nextValue)} - /> + onKeyValueChange("basic_info", nextValue)} /> + onKeyValueChange("emotional_patterns", nextValue)} /> + onKeyValueChange("preferences", nextValue)} />
- +
@@ -99,9 +70,7 @@ export function MemoryDesk(props: MemoryDeskProps) {

{conversation.agent_message}

))} - {!draft.recent_conversations?.length ? ( -

这个用户还没有历史对话,保存记忆后可直接用于回复预览。

- ) : null} + {!draft.recent_conversations?.length ?

这个用户还没有历史对话。

: null}
diff --git a/admin-ui/src/components/ProactiveStudio.tsx b/admin-ui/src/components/ProactiveStudio.tsx index 0b453fc..b9edb2b 100644 --- a/admin-ui/src/components/ProactiveStudio.tsx +++ b/admin-ui/src/components/ProactiveStudio.tsx @@ -1,4 +1,4 @@ -import type { ProactiveChatConfig, UserSummary } from "../types"; +import type { ProactiveChatConfig, UserSummary } from "../types"; type ProactiveStudioProps = { config: ProactiveChatConfig; @@ -9,7 +9,7 @@ type ProactiveStudioProps = { previewReply: string; deliveryStatus: string; onToggleEnabled: (value: boolean) => void; - onTargetUserChange: (value: string) => void; + onTargetUserChange: (channel: string, externalUserId: string) => void; onWindowToggle: (key: string, enabled: boolean) => void; onWindowTimeChange: (key: string, value: string) => void; onQuietHoursToggle: (value: boolean) => void; @@ -43,6 +43,8 @@ export function ProactiveStudio(props: ProactiveStudioProps) { onRunOnce, } = props; + const targetValue = config.target_external_user_id ? `${config.target_channel}::${config.target_external_user_id}` : ""; + return (
@@ -64,11 +66,22 @@ export function ProactiveStudio(props: ProactiveStudioProps) {
diff --git a/admin-ui/src/components/ProactiveStudio.tsx b/admin-ui/src/components/ProactiveStudio.tsx index 0b453fc..b9edb2b 100644 --- a/admin-ui/src/components/ProactiveStudio.tsx +++ b/admin-ui/src/components/ProactiveStudio.tsx @@ -1,4 +1,4 @@ -import type { ProactiveChatConfig, UserSummary } from "../types"; +import type { ProactiveChatConfig, UserSummary } from "../types"; type ProactiveStudioProps = { config: ProactiveChatConfig; @@ -9,7 +9,7 @@ type ProactiveStudioProps = { previewReply: string; deliveryStatus: string; onToggleEnabled: (value: boolean) => void; - onTargetUserChange: (value: string) => void; + onTargetUserChange: (channel: string, externalUserId: string) => void; onWindowToggle: (key: string, enabled: boolean) => void; onWindowTimeChange: (key: string, value: string) => void; onQuietHoursToggle: (value: boolean) => void; @@ -43,6 +43,8 @@ export function ProactiveStudio(props: ProactiveStudioProps) { onRunOnce, } = props; + const targetValue = config.target_external_user_id ? `${config.target_channel}::${config.target_external_user_id}` : ""; + return (
@@ -64,11 +66,22 @@ export function ProactiveStudio(props: ProactiveStudioProps) {