# agent_chat.dspy - 通用会话 Agent(pipeline-core,产线无关) # 使用 AgentExecutor v2(gateway 统一入口),是 Web 端复刻 Hermes CLI 的通用交互能力 # 各产线/模块可在其上叠加专属能力,本 dspy 不依赖任何产线专属表 import json import os import re import zipfile from ahserver.filestorage import FileStorage def _extract_text(path, name): """提取文件文本内容(docx/txt/md等),返回文本或空字符串。二进制/无法解析返回空。""" ext = os.path.splitext(name)[1].lower() try: if ext in ('.txt', '.md', '.json', '.csv', '.py', '.log', '.yaml', '.yml', '.xml', '.html', '.ini'): with open(path, 'r', encoding='utf-8', errors='ignore') as f: return f.read()[:15000] if ext == '.docx': with zipfile.ZipFile(path) as z: xml = z.read('word/document.xml').decode('utf-8', errors='ignore') texts = re.findall(r']*>(.*?)', xml) return '\n'.join(texts)[:15000] except Exception: pass return '' action = (params_kw or {}).get('action', 'send_message') msg = '' if action != 'send_message': p = params_kw or {} msg = (p.get('message_text') or '').strip() if not msg: inner = p.get('params', {}) if isinstance(inner, dict): msg = (inner.get('prompt') or inner.get('message_text') or '').strip() if msg: action = 'send_message' dbname = get_module_dbname('pipeline_core') if action == 'send_message': # 从 params 中提取 prompt prompt = (params_kw or {}).get('prompt', '') or msg prompt = prompt.strip() # 处理用户上传的文件(multipart file 字段 → web_path),提取文本作为中性上下文注入 prompt file_ctx = '' _fval = (params_kw or {}).get('file') _fpaths = _fval if isinstance(_fval, list) else ([_fval] if _fval else []) for _fp in _fpaths: try: _abs = FileStorage().realPath(_fp) _name = os.path.basename(_abs) _txt = _extract_text(_abs, _name) if _txt: file_ctx += f"【文件 {_name} 内容】\n{_txt}\n\n" else: file_ctx += f"【文件 {_name}】二进制文件,无法直接读取文本。\n\n" except Exception: pass if file_ctx: prompt = "用户本次上传了以下文件:\n\n" + file_ctx + "\n用户指令:" + prompt # ── 意图识别:停止类指令 ── stop_keywords = ['停止', '取消', '停下', '终止', '停', 'stop', 'cancel', 'abort'] if any(kw == prompt.strip().lower() or kw in prompt.strip().lower() for kw in stop_keywords): async def _stop_stream(): yield json.dumps({"content": "已停止当前任务。"}, ensure_ascii=False) + '\n' return await stream_response(request, _stop_stream, 'text/plain; charset=utf-8') uid = await get_user() if not uid: uid = 'user-01' # 测试兼容 # ── 走 gateway 统一入口(Web AgentIO 通道)── from pipeline_service.gateway import get_gateway gateway = get_gateway() async def agent_stream(): async for chunk in gateway.run_message("web", uid, prompt): data = json.loads(chunk) t = data.get('type', '') if t == 'progress': yield json.dumps({"reasoning_content": data.get('message', '') + "\n"}, ensure_ascii=False) + '\n' elif t == 'auto_tool': yield json.dumps({"reasoning_content": "🔄 " + data.get('message', '') + "\n"}, ensure_ascii=False) + '\n' elif t == 'debug': continue # skip debug in UI elif t == 'tool_call': tool = data.get('tool', '') params = json.dumps(data.get('params', {}), ensure_ascii=False) yield json.dumps({"content": "**🔧 调用: " + tool + "**\n```\n" + params + "\n```\n\n"}, ensure_ascii=False) + '\n' elif t == 'tool_result': result = data.get('result', '') yield json.dumps({"content": result + "\n\n"}, ensure_ascii=False) + '\n' elif t == 'reply': yield json.dumps({"content": data.get('message', '')}, ensure_ascii=False) + '\n' elif t == 'ask_user': yield json.dumps({"content": "❓ " + data.get('message', '')}, ensure_ascii=False) + '\n' elif t == 'error': yield json.dumps({"error": data.get('message', '')}, ensure_ascii=False) + '\n' else: yield json.dumps({"content": chunk}, ensure_ascii=False) + '\n' return await stream_response(request, agent_stream, 'text/plain; charset=utf-8') elif action == 'list_messages': uid = await get_user() if not uid: return json.dumps({"error": "请先登录"}, ensure_ascii=False) async with DBPools().sqlorContext(dbname) as sor: ms = await sor.sqlExe( "SELECT role,content,created_at FROM pipeline_conversations ORDER BY created_at ASC LIMIT 100", {}) result = [{"role": getattr(m, 'role', ''), "content": getattr(m, 'content', '')} for m in (ms or [])] return json.dumps({"success": True, "messages": result}, ensure_ascii=False, default=str) else: return json.dumps({"error": "Unknown: " + str(action)}, ensure_ascii=False)