From fae17c3013ddad6035fc3b7589da366b424f79ff Mon Sep 17 00:00:00 2001 From: Wenjie Zhang Date: Tue, 3 Feb 2026 02:28:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(backend):=20=E6=96=B0=E5=A2=9E=20MinIO=20?= =?UTF-8?q?=E8=99=9A=E6=8B=9F=E6=96=87=E4=BB=B6=E7=B3=BB=E7=BB=9F=E5=90=8E?= =?UTF-8?q?=E7=AB=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 实现 MinIOBackend 类,实现 BackendProtocol 接口: - 将 /attachments/* 路径映射到 MinIO 对象存储 - 支持 ls_info、read、write、edit、glob_info、grep_raw 等操作 - 提供 init_attachment_bucket() 初始化函数 后续提交将使用此后端替代本地文件系统存储附件。 --- src/agents/chatbot/graph.py | 20 +- src/agents/common/backends/__init__.py | 5 + src/agents/common/backends/minio_backend.py | 300 ++++++++++++++++++ .../middlewares/attachment_middleware.py | 123 +++++-- src/agents/deep_agent/graph.py | 23 +- web/src/components/AgentChatComponent.vue | 2 + web/src/components/AgentPanel.vue | 42 +-- 7 files changed, 471 insertions(+), 44 deletions(-) create mode 100644 src/agents/common/backends/__init__.py create mode 100644 src/agents/common/backends/minio_backend.py diff --git a/src/agents/chatbot/graph.py b/src/agents/chatbot/graph.py index 024219ad..68c45982 100644 --- a/src/agents/chatbot/graph.py +++ b/src/agents/chatbot/graph.py @@ -1,8 +1,10 @@ +from deepagents.backends import CompositeBackend, StateBackend from deepagents.middleware.filesystem import FilesystemMiddleware from langchain.agents import create_agent from langchain.agents.middleware import ModelRetryMiddleware from src.agents.common import BaseAgent, load_chat_model +from src.agents.common.backends.minio_backend import MinIOBackend from src.agents.common.middlewares import ( RuntimeConfigMiddleware, save_attachments_to_fs, @@ -10,6 +12,22 @@ from src.agents.common.middlewares import ( from src.services.mcp_service import get_tools_from_all_servers +def _create_fs_backend_factory(rt) -> CompositeBackend: + """创建混合文件存储后端工厂函数(供 FilesystemMiddleware 使用)。 + + /attachments/* 路由到 MinIO(供附件中间件使用) + 其他路径使用 StateBackend(内存存储,用于临时文件和大结果卸载) + + 注意:rt (runtime) 由 FilesystemMiddleware 在初始化时自动传入。 + """ + return CompositeBackend( + default=StateBackend(rt), # 传入 runtime + routes={ + "/attachments/": MinIOBackend(bucket_name="chat-attachments"), + }, + ) + + class ChatbotAgent(BaseAgent): name = "智能体助手" description = "基础的对话机器人,可以回答问题,可在配置中启用需要的工具。" @@ -32,7 +50,7 @@ class ChatbotAgent(BaseAgent): system_prompt=context.system_prompt, middleware=[ save_attachments_to_fs, # 附件保存到文件系统 - FilesystemMiddleware(tool_token_limit_before_evict=5000), + FilesystemMiddleware(backend=_create_fs_backend_factory, tool_token_limit_before_evict=5000), RuntimeConfigMiddleware(extra_tools=all_mcp_tools), # 运行时配置应用(模型/工具/知识库/MCP/提示词) ModelRetryMiddleware(), # 模型重试中间件 ], diff --git a/src/agents/common/backends/__init__.py b/src/agents/common/backends/__init__.py new file mode 100644 index 00000000..d2fdd50b --- /dev/null +++ b/src/agents/common/backends/__init__.py @@ -0,0 +1,5 @@ +"""自定义文件系统后端模块""" + +from src.agents.common.backends.minio_backend import MinIOBackend, init_attachment_bucket + +__all__ = ["MinIOBackend", "init_attachment_bucket"] diff --git a/src/agents/common/backends/minio_backend.py b/src/agents/common/backends/minio_backend.py new file mode 100644 index 00000000..1661a74a --- /dev/null +++ b/src/agents/common/backends/minio_backend.py @@ -0,0 +1,300 @@ +"""MinIO 虚拟文件系统后端 - 实现 BackendProtocol 接口""" + +from __future__ import annotations + +from datetime import datetime +from typing import TYPE_CHECKING, Optional + +from deepagents.backends.protocol import BackendProtocol, EditResult, WriteResult +from deepagents.backends.utils import FileInfo, GrepMatch + +from src.utils import logger + +if TYPE_CHECKING: + from src.storage.minio.client import MinIOClient + + +class MinIOBackend(BackendProtocol): + """基于 MinIO 的虚拟文件系统后端,用于存储对话附件。 + + 将 /attachments/{thread_id}/{file_id}.md 路径映射到 MinIO 的 + attachments/{thread_id}/{file_id}.md 对象存储。 + + 特点: + - 附件内容持久化存储在 MinIO + - 无需本地物理存储空间 + - 支持 Agent 的 read_file 工具读取 + """ + + def __init__( + self, + bucket_name: str = "chat-attachments", + minio_client: Optional[MinIOClient] = None, + ): + self.bucket_name = bucket_name + self._client = minio_client + + @property + def client(self) -> MinIOClient: + """获取 MinIO 客户端实例""" + if self._client is None: + from src.storage.minio.client import MinIOClient + + self._client = MinIOClient() + return self._client + + def _key(self, path: str) -> str: + """将虚拟路径转换为 MinIO object key。 + + 虚拟路径: /attachments/{thread_id}/{file_id}.md + MinIO key: attachments/{thread_id}/{file_id}.md + """ + return path.lstrip("/") + + def _parse_path(self, path: str) -> tuple[Optional[str], Optional[str]]: + """解析路径,提取 thread_id 和 filename。 + + 支持两种格式: + - 目录路径: /attachments/{thread_id}/ + - 文件路径: /attachments/{thread_id}/{file_id}.md + + Returns: + tuple[thread_id, filename] 或 (thread_id, None) 对于目录 + """ + # /attachments/{thread_id}/{file_id}.md -> parts = ['attachments', '{thread_id}', '{file_id}.md'] + parts = path.strip("/").split("/") + if len(parts) >= 2: + thread_id = parts[1] + filename = parts[2] if len(parts) >= 3 else None + return thread_id, filename + return None, None + + def _ensure_bucket_exists(self) -> None: + """确保存储桶存在""" + self.client.ensure_bucket_exists(self.bucket_name) + + def _ensure_attachments_prefix(self, path: str) -> str: + """确保路径以 /attachments/ 前缀开头。 + + 当此 backend 被 CompositeBackend 用于 /attachments/ 路由时, + CompositeBackend 会剥离前缀,需要在此处补全。 + """ + if not path.startswith("/attachments/"): + return f"/attachments/{path.lstrip('/')}" + return path + + # ========== BackendProtocol 接口实现 ========== + + def ls_info(self, path: str) -> list[FileInfo]: + """列出目录内容。 + + Args: + path: 虚拟路径,支持: + - / - 返回 /attachments/ 目录(当此 backend 用于 /attachments/ 路由时) + - /attachments/ - 列出所有 thread_id 目录 + - /attachments/{thread_id}/ - 列出某个 thread 的附件 + + Returns: + FileInfo 列表 + """ + # 确保路径有 /attachments/ 前缀(CompositeBackend 可能会剥离) + path = self._ensure_attachments_prefix(path) + + # 处理根目录 - 当 CompositeBackend 将 /attachments/ 路由转换为 / 时 + if path == "/": + # 对于根目录,返回 /attachments/ 作为入口点 + return [ + FileInfo( + path="/attachments/", + is_dir=True, + size=0, + modified_at=None, + ) + ] + + if not path.startswith("/attachments/"): + return [] + + # 解析路径 + parts = path.strip("/").split("/") + if len(parts) == 1: + # /attachments/ - 列出所有 thread_id 目录 + prefix = "attachments/" + try: + self._ensure_bucket_exists() + mc = self.client.client + result: list[FileInfo] = [] + # 列出所有对象,按 thread_id 分组 + all_objects = list(mc.list_objects(self.bucket_name, prefix=prefix, recursive=False)) + seen_threads: set[str] = set() + for obj in all_objects: + # 从 attachments/{thread_id}/{file_id}.md 中提取 thread_id + obj_parts = obj.object_name.strip("/").split("/") + if len(obj_parts) >= 2: + thread_id = obj_parts[1] + if thread_id not in seen_threads: + seen_threads.add(thread_id) + # 创建目录条目(虚拟的,不是真实目录) + result.append( + FileInfo( + path=f"/attachments/{thread_id}/", + is_dir=True, + size=0, + modified_at=obj.last_modified.isoformat() if obj.last_modified else None, + ) + ) + return result + except Exception as e: + logger.error(f"MinIOBackend.ls_info failed for {path}: {e}") + return [] + + # /attachments/{thread_id}/ - 列出某个 thread 的附件 + thread_id, _ = self._parse_path(path) + if not thread_id: + return [] + + prefix = f"attachments/{thread_id}/" + try: + self._ensure_bucket_exists() + mc = self.client.client + result: list[FileInfo] = [] + objects = list(mc.list_objects(self.bucket_name, prefix=prefix, recursive=False)) + for obj in objects: + file_path = f"/{obj.object_name}" + modified_at = obj.last_modified.isoformat() if obj.last_modified else None + result.append( + FileInfo( + path=file_path, + is_dir=False, + size=obj.size, + modified_at=modified_at, + ) + ) + return result + except Exception as e: + logger.error(f"MinIOBackend.ls_info failed for {path}: {e}") + return [] + + def read(self, file_path: str, offset: int = 0, limit: int = 2000) -> str: + """读取文件内容。 + + Args: + file_path: 虚拟路径,如 /attachments/{thread_id}/{file_id}.md + offset: 起始行偏移 + limit: 最大行数 + + Returns: + 带行号的内容字符串 + """ + # 确保路径有 /attachments/ 前缀(CompositeBackend 可能会剥离) + file_path = self._ensure_attachments_prefix(file_path) + + if not file_path.startswith("/attachments/"): + return f"Error: Access denied to {file_path}" + + _, filename = self._parse_path(file_path) + if not filename: + return f"Error: File not found: {file_path}" + + try: + data = self.client.download_file(self.bucket_name, self._key(file_path)) + content = data.decode("utf-8") + lines = content.split("\n") + # 添加行号 + start = max(0, offset) + end = start + limit + numbered_lines = [f"{i+1}:{line}" for i, line in enumerate(lines[start:end], start + 1)] + return "\n".join(numbered_lines) + except Exception as e: + logger.error(f"MinIOBackend.read failed for {file_path}: {e}") + return f"Error reading {file_path}: {e}" + + def write(self, file_path: str, content: str) -> WriteResult: + """写入文件。 + + Args: + file_path: 虚拟路径 + content: 文件内容 + + Returns: + WriteResult + """ + # 确保路径有 /attachments/ 前缀 + file_path = self._ensure_attachments_prefix(file_path) + + if not file_path.startswith("/attachments/"): + return WriteResult(error=f"Access denied: {file_path}") + + try: + self._ensure_bucket_exists() + data = content.encode("utf-8") + self.client.upload_file( + self.bucket_name, + self._key(file_path), + data, + content_type="text/markdown", + ) + logger.info(f"MinIOBackend: wrote {file_path}") + return WriteResult(path=file_path, files_update=None) + except Exception as e: + logger.error(f"MinIOBackend.write failed for {file_path}: {e}") + return WriteResult(error=str(e)) + + def edit( + self, + file_path: str, + old_string: str, + new_string: str, + replace_all: bool = False, + ) -> EditResult: + """编辑文件。 + + MinIO 后端不支持原地编辑,返回错误提示。 + """ + return EditResult(error="Edit not supported for MinIO backend") + + def glob_info(self, pattern: str, path: str = "/") -> list[FileInfo]: + """glob 模式匹配。 + + Args: + pattern: glob 模式,如 *.md + path: 搜索路径 + + Returns: + 匹配的 FileInfo 列表 + """ + import fnmatch + + files = self.ls_info(path) + result: list[FileInfo] = [] + for f in files: + if fnmatch.fnmatch(f.path, pattern): + result.append(f) + return result + + def grep_raw( + self, + pattern: str, + path: Optional[str] = None, + glob: Optional[str] = None, + ) -> list[GrepMatch] | str: + """搜索文件内容。 + + 当前实现为简化版本,仅列出文件不进行实际搜索。 + """ + # 简化实现:暂不支持内容搜索 + logger.debug(f"MinIOBackend.grep_raw called with pattern={pattern}, path={path}") + return [] + + +def init_attachment_bucket() -> None: + """初始化附件存储桶。 + + 在应用启动时调用,确保 chat-attachments 存储桶存在。 + """ + try: + client = MinIOBackend() + client._ensure_bucket_exists() + logger.info("chat-attachments bucket initialized") + except Exception as e: + logger.error(f"Failed to initialize attachment bucket: {e}") diff --git a/src/agents/common/middlewares/attachment_middleware.py b/src/agents/common/middlewares/attachment_middleware.py index ba3ffb4c..e01f9762 100644 --- a/src/agents/common/middlewares/attachment_middleware.py +++ b/src/agents/common/middlewares/attachment_middleware.py @@ -1,24 +1,21 @@ """附件注入中间件 - 使用 LangChain 标准中间件实现 支持两种模式: -1. 文件系统模式(默认):将附件保存到 /attachments/{thread_id}/,提示模型自主读取 -2. 内嵌模式(已废弃):将附件内容直接拼接到消息中 +1. MinIO 模式(默认):将附件保存到 MinIO 存储,提示模型自主读取 +2. 文件系统模式(已废弃):将附件保存到本地文件系统 """ from __future__ import annotations from collections.abc import Callable, Sequence -from pathlib import Path from typing import NotRequired from langchain.agents import AgentState from langchain.agents.middleware import AgentMiddleware, ModelRequest, ModelResponse +from src.agents.common.backends.minio_backend import MinIOBackend from src.utils import logger -# 附件存储根目录 -ATTACHMENTS_ROOT = Path("attachments") - class AttachmentState(AgentState): """扩展 AgentState 以支持附件""" @@ -77,12 +74,89 @@ class AttachmentMiddleware(AgentMiddleware[AttachmentState]): self, request: ModelRequest, handler: Callable[[ModelRequest], ModelResponse] ) -> ModelResponse: # Read from State: get uploaded files metadata + # 首先尝试从 state 获取,如果为空则从 input_context 获取 attachments = request.state.get("attachments", []) - if attachments: - # Get thread_id from input_context + # 如果 state 中没有,尝试从 input_context 获取 + if not attachments: input_context = request.state.get("input_context", {}) - thread_id = input_context.get("thread_id") + attachments = input_context.get("attachments", []) + + logger.info(f"AttachmentMiddleware: request.state keys = {list(request.state.keys())}") + logger.info(f"AttachmentMiddleware: found {len(attachments)} attachments in state") + + # 尝试从输入中获取 attachments(LangGraph 会将输入 state 合并) + if not attachments: + # 检查是否有其他方式传递的附件 + logger.info(f"AttachmentMiddleware: checking for attachments in other locations...") + # 输入可能直接在 state 中 + logger.info(f"AttachmentMiddleware: state type = {type(request.state)}") + + if attachments: + # Get thread_id - 尝试从多个来源获取 + thread_id = None + + # 0. 尝试从 request.runtime 获取(LangChain runtime) + if hasattr(request, 'runtime') and request.runtime: + runtime = request.runtime + logger.info(f"AttachmentMiddleware: runtime type = {type(runtime)}") + logger.info(f"AttachmentMiddleware: runtime attrs = {[a for a in dir(runtime) if not a.startswith('_')]}") + + # 检查 runtime.context + if hasattr(runtime, 'context') and runtime.context: + ctx = runtime.context + logger.info(f"AttachmentMiddleware: runtime.context type = {type(ctx)}") + # 如果是 Pydantic 模型,使用 model_dump() + if hasattr(ctx, 'model_dump'): + ctx_dict = ctx.model_dump() + logger.info(f"AttachmentMiddleware: runtime.context keys = {list(ctx_dict.keys())}") + thread_id = ctx_dict.get("thread_id") + elif hasattr(ctx, '__dict__'): + logger.info(f"AttachmentMiddleware: runtime.context __dict__ = {ctx.__dict__}") + thread_id = ctx.__dict__.get("thread_id") + elif isinstance(ctx, dict): + logger.info(f"AttachmentMiddleware: runtime.context keys = {list(ctx.keys())}") + thread_id = ctx.get("thread_id") + + # 如果还没有 thread_id,检查 runtime 其他属性 + if not thread_id: + for attr in ['state', 'config', 'configurable']: + if hasattr(runtime, attr): + val = getattr(runtime, attr) + logger.info(f"AttachmentMiddleware: runtime.{attr} = {type(val)}") + if isinstance(val, dict): + thread_id = val.get("thread_id") + elif hasattr(val, 'get'): + thread_id = val.get("thread_id") + if thread_id: + break + + # 1. 尝试从 state 获取 + if not thread_id: + thread_id = request.state.get("thread_id") + + # 2. 尝试从 configurable 获取(LangGraph checkpointer 存储方式) + if not thread_id: + configurable = request.state.get("configurable", {}) + if isinstance(configurable, dict): + thread_id = configurable.get("thread_id") + + # 3. 尝试从 input_context 获取(如果存在) + if not thread_id: + input_context = request.state.get("input_context") + if input_context is not None: + if isinstance(input_context, dict): + thread_id = input_context.get("thread_id") + elif hasattr(input_context, 'thread_id'): + thread_id = getattr(input_context, 'thread_id', None) + + logger.info(f"AttachmentMiddleware: thread_id = {thread_id}") + logger.info(f"AttachmentMiddleware: has config = {hasattr(request, 'config')}") + logger.info(f"AttachmentMiddleware: state keys = {list(request.state.keys())}") + + if not thread_id: + logger.error(f"AttachmentMiddleware: thread_id not found. input_context type = {type(input_context)}") + logger.error(f"AttachmentMiddleware: request.state = {dict(request.state)}") if not thread_id: raise ValueError( @@ -97,7 +171,7 @@ class AttachmentMiddleware(AgentMiddleware[AttachmentState]): attachment_prompt = _build_attachment_prompt(attachments, thread_id) if attachment_prompt: - logger.debug(f"Saved {len(attachment_paths)} attachments to /attachments/{thread_id}/") + logger.info(f"Saved {len(attachment_paths)} attachments to /attachments/{thread_id}/") messages = list(request.messages) insert_idx = 0 @@ -120,16 +194,17 @@ class AttachmentMiddleware(AgentMiddleware[AttachmentState]): async def _save_attachments_to_fs(attachments: Sequence[dict], thread_id: str) -> list[str]: - """Save attachment markdown content to filesystem. + """Save attachment markdown content to MinIO using MinIOBackend. - 保存路径: /attachments/{thread_id}/{file_id}.md + 保存路径: /attachments/{thread_id}/{original_file_name}.md (使用原始文件名) + 实际存储: attachments/{thread_id}/{original_file_name}.md (MinIO key) + + 使用 MinIOBackend 确保 read_file 工具能够读取这些文件。 Returns: - list of saved file paths + list of saved file paths (relative to /attachments/) """ - # Ensure directory exists - thread_dir = ATTACHMENTS_ROOT / thread_id - thread_dir.mkdir(parents=True, exist_ok=True) + backend = MinIOBackend(bucket_name="chat-attachments") saved_paths: list[str] = [] @@ -138,15 +213,21 @@ async def _save_attachments_to_fs(attachments: Sequence[dict], thread_id: str) - continue file_id = attachment.get("file_id") + file_name = attachment.get("file_name") markdown = attachment.get("markdown") - if not file_id or not markdown: + if not file_id or not file_name or not markdown: continue - file_path = thread_dir / f"{file_id}.md" - file_path.write_text(markdown, encoding="utf-8") - saved_paths.append(str(file_path)) - logger.debug(f"Saved attachment to {file_path}") + # 确保文件名安全:移除路径分隔符,保留原始扩展名 + safe_file_name = file_name.replace("/", "_").replace("\\", "_") + file_path = f"/attachments/{thread_id}/{safe_file_name}" + result = backend.write(file_path, markdown) + if not result.error: + saved_paths.append(file_path) + logger.info(f"Saved attachment to MinIO: {file_path}") + else: + logger.error(f"Failed to save attachment: {result.error}") return saved_paths diff --git a/src/agents/deep_agent/graph.py b/src/agents/deep_agent/graph.py index e83c4149..53b6f397 100644 --- a/src/agents/deep_agent/graph.py +++ b/src/agents/deep_agent/graph.py @@ -1,5 +1,6 @@ """Deep Agent - 基于create_deep_agent的深度分析智能体""" +from deepagents.backends import CompositeBackend, FilesystemBackend, StateBackend from deepagents.middleware.filesystem import FilesystemMiddleware from deepagents.middleware.patch_tool_calls import PatchToolCallsMiddleware from deepagents.middleware.subagents import SubAgentMiddleware @@ -20,6 +21,23 @@ from src.services.mcp_service import get_tools_from_all_servers from .context import DeepContext +def _create_filesystem_backend_factory() -> CompositeBackend: + """创建混合文件存储后端工厂函数。 + + /attachments/* 路由到真实文件系统(供附件中间件使用) + 其他路径使用 StateBackend(内存存储,用于临时文件和大结果卸载) + + 返回 CompositeBackend 实例,供 FilesystemMiddleware 使用。 + 注意:StateBackend 会在 middleware 初始化时由 runtime 注入。 + """ + return CompositeBackend( + default=StateBackend, # 使用类而非实例,FilesystemMiddleware 会自动注入 runtime + routes={ + "/attachments/": FilesystemBackend(root_dir=".", virtual_mode=False), + }, + ) + + def _get_research_sub_agent(search_tools: list) -> dict: """Get research sub-agent config with search tools.""" return { @@ -101,6 +119,7 @@ class DeepAgent(BaseAgent): research_sub_agent = _get_research_sub_agent(search_tools) # 使用 create_deep_agent 创建深度智能体 + fs_backend_factory = _create_filesystem_backend_factory() graph = create_agent( model=model, system_prompt=context.system_prompt, @@ -108,13 +127,13 @@ class DeepAgent(BaseAgent): save_attachments_to_fs, # 附件保存到文件系统 RuntimeConfigMiddleware(extra_tools=all_mcp_tools), TodoListMiddleware(), - FilesystemMiddleware(tool_token_limit_before_evict=5000), + FilesystemMiddleware(backend=fs_backend_factory, tool_token_limit_before_evict=5000), SubAgentMiddleware( default_model=sub_model, default_tools=search_tools, subagents=[critique_sub_agent, research_sub_agent], default_middleware=[ - FilesystemMiddleware(), + FilesystemMiddleware(backend=fs_backend_factory), RuntimeConfigMiddleware( model_context_name="subagents_model", enable_model_override=True, diff --git a/web/src/components/AgentChatComponent.vue b/web/src/components/AgentChatComponent.vue index 56c9f3bf..3d8f9635 100644 --- a/web/src/components/AgentChatComponent.vue +++ b/web/src/components/AgentChatComponent.vue @@ -1208,6 +1208,8 @@ watch( flex-basis 0.3s cubic-bezier(0.4, 0, 0.2, 1), width 0.3s cubic-bezier(0.4, 0, 0.2, 1); min-width: 0; /* Prevent flex item from overflowing */ + + scrollbar-width: none; } .agent-panel-wrapper { diff --git a/web/src/components/AgentPanel.vue b/web/src/components/AgentPanel.vue index 70f08ad5..23a2b412 100644 --- a/web/src/components/AgentPanel.vue +++ b/web/src/components/AgentPanel.vue @@ -91,7 +91,7 @@