# 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): """[已废弃] 保留签名兼容;真实逻辑在 pipeline_core.upload_tools.extract_text。""" from pipeline_core.upload_tools import extract_text preview, _total, _trunc = extract_text(path, name) return preview 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): # 统一走 pipeline_core.upload_tools——截断显式告知、落盘位置=agent 可读位置。 from pipeline_core.upload_tools import extract_text, resolve_upload_dir, save_uploads, build_file_context _items = [] # build_file_context 入参 _uploads = [] # [(src_abs, filename)] _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) _uploads.append((_abs, _name)) except Exception: pass # ── 意图识别:停止类指令 ── # 注意(2026-08-31 修复):会话 agent 新增了项目管理能力(终止/删除项目走 # pause_project/delete_project 工具,代码层校验人类 owner)。「终止项目」「停止项目」 # 含停止词但语义是项目管理指令,不能拦——只对纯停止指令(或不含项目词的)拦截。 stop_keywords = ['停止', '取消', '停下', '终止', '停', 'stop', 'cancel', 'abort'] _pl = prompt.strip().lower() _is_project_mgmt = any(w in _pl for w in ('项目', '工程', 'project')) if any(kw == _pl or (kw in _pl and not _is_project_mgmt) 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' # 测试兼容 # ── 会话内唯一标识(web 多 tab 独立会话):前端每个 tab 传独立 session_id ── session_id = (params_kw or {}).get('session_id', '') or '' # ── 产线默认能力(入口指定,如 bidding_general):无当前项目时装载该产线能力 ── pipeline_id = (params_kw or {}).get('pipeline_id', '') or '' # ── 上传文件落盘(有项目→项目根;无项目→workspace 会话目录)+ 生成显式截断上下文 ── file_ctx = '' if _uploads: try: async with DBPools().sqlorContext(dbname) as sor: _udir, _rel = await resolve_upload_dir(sor, uid, session_id) _saved = save_uploads(_udir, _uploads) _byname = {os.path.basename(src): n for (src, _n), (n, _p) in zip(_uploads, _saved)} for _src, _name in _uploads: _preview, _total, _trunc = extract_text(_src, _name) _isbin = (not _preview and _total == 0) _relp = (_rel + _byname.get(_name, _name)) if _name in _byname else '' _items.append((_name, _relp, _preview, _total, _trunc, _isbin)) file_ctx = build_file_context(_items) except Exception: file_ctx = '' if file_ctx: prompt = file_ctx + "\n用户指令:" + prompt # ── 走 gateway 统一入口(Web AgentIO 通道)── # model_id:前端模型下拉选中的模型 → gateway 校验后持久化到项目(项目模型一经设置 # 即生效,直到用户再次选择),空 = 沿用项目已设模型。 model_id = (params_kw or {}).get('model_id', '') or '' from pipeline_service.gateway import get_gateway gateway = get_gateway() async def agent_stream(): _base = entire_url("/") async for chunk in gateway.run_message("web", uid, prompt, base_url=_base, session_id=session_id, pipeline_id=pipeline_id, model_id=model_id): 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': msg = data.get('message', '') if isinstance(msg, dict) and msg.get('widgettype'): # widget JSON:顶层透传,前端 AgentOut 检测 widgettype 渲染(如 /task 的 PopupWindow) yield json.dumps(msg, ensure_ascii=False) + '\n' else: yield json.dumps({"content": msg}, ensure_ascii=False) + '\n' elif t == 'ask_user': yield json.dumps({"content": "❓ " + data.get('message', '')}, ensure_ascii=False) + '\n' elif t == 'confirm': _tool = data.get('tool', '') _params = data.get('params', {}) or {} _cmd = _params.get('command', '') if isinstance(_params, dict) else '' yield json.dumps({"content": "⚠️ 需要你确认执行命令:`" + str(_cmd)[:120] + "`\n请回复「确认」执行、「全部确认」本会话不再询问,或「取消」放弃。"}, 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)