ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Cua MCP Server 并发会话管理:SessionManager 与 ComputerPool 的源码级解析

Cua MCP Server 并发会话管理:SessionManager 与 ComputerPool 的源码级解析 Cua MCP Server 并发会话管理SessionManager 与 ComputerPool 的源码级解析【免费下载链接】cuaScale computer-use 2.0 with open-source drivers, cross-OS fleets, and benchmarks for training, evaluation, and data generation.项目地址: https://gitcode.com/GitHub_Trending/cua/cuaMCPModel Context ProtocolServer 让 Cua 的 Computer-Use Agent 能够接入 Claude Desktop、Cursor 等 MCP 客户端但当多个客户端同时连接时早期实现中的全局单例 Computer 实例、任务串行处理、缺乏资源回收等问题会直接破坏多客户端体验。本文基于仓库文档 CONCURRENT_SESSIONS.md结合 session_manager.py 与 server.py 的实际源码完整拆解这套并发会话管理机制每个客户端如何获得隔离的 Computer 实例、实例池如何复用资源、任务如何并发执行、服务如何优雅关闭并给出可复制的多客户端接入示例与测试验证路径。1. 问题陈述旧实现为什么无法支撑并发原始 MCP Server 实现存在五个关键问题引自 CONCURRENT_SESSIONS.md 的 Problem Statement全局 Computer 实例所有客户端共享一个global_computer变量无资源隔离多个客户端会互相干扰任务串行处理多任务操作只能顺序执行无优雅关闭服务关闭时无法正确清理资源隐藏的事件循环server.run()隐藏了事件循环无法进行正确的生命周期管理。在 server.py 中可以看到问题 5 的修复痕迹——注释明确写道“Use run_stdio_async directly instead of server.run() to avoid nested event loops”即改用server.run_stdio_async()直接驱动 stdio 传输把事件循环的控制权交还给run_server()从而使信号处理和清理逻辑得以在同一事件循环内运行。2. 核心架构SessionManager 与 ComputerPool解决方案位于 libs/python/mcp-server/mcp_server/session_manager.py由三层构成SessionInfo数据类、ComputerPool实例池、SessionManager会话管理器。2.1 SessionInfo会话状态的数据结构每个会话用一个 dataclass 承载全部生命周期状态session_manager.pydataclass class SessionInfo: Information about an active session. session_id: str computer: Any # Computer instance created_at: float last_activity: float active_tasks: Set[str] field(default_factoryset) is_shutting_down: bool False字段语义直接对应文档中的会话生命周期四阶段创建 → 任务注册 → 活动追踪 → 清理last_activity用于空闲判定active_tasks防止清理有活跃任务的会话is_shutting_down则拦截对正在清理的会话的新请求。2.2 ComputerPoolComputer 实例池ComputerPool负责 Computer 实例的复用与回收默认参数为max_size: int 5, idle_timeout: float 300.0session_manager.pyclass ComputerPool: Pool of computer instances for efficient resource management. def __init__(self, max_size: int 5, idle_timeout: float 300.0): self.max_size max_size self.idle_timeout idle_timeout self._available: List[Any] [] self._in_use: Set[Any] set() self._creation_lock asyncio.Lock()其acquire()方法体现了三段式获取策略session_manager.py复用空闲实例若_available队列非空直接弹出并计入_in_use避免重复启动 VM 的开销受锁保护地创建新实例在_creation_lock临界区内检查len(self._in_use) self.max_size通过后创建Computer并await computer.run()完成启动。源码中还包含一个值得注意的配置开关use_host os.getenv(CUA_USE_HOST_COMPUTER_SERVER, false).lower() in ( true, 1, yes, ) computer Computer(verbositylogging.INFO, use_host_computer_serveruse_host)即通过环境变量CUA_USE_HOST_COMPUTER_SERVER可让池中的实例连接到宿主机上已有的 computer-server而不是由池自行启动这是一个文档未提及但从源码可直接确认的部署选项 3.轮询等待池满时进入while not self._available: await asyncio.sleep(0.1)的等待循环直到有实例被释放回来。release()把实例从_in_use移回_availableshutdown()则对两个集合中的实例统一调用close()不存在时回退到stop()并清空集合保证关闭路径不泄漏任何 VM。2.3 SessionManager会话编排与自动清理SessionManager对外暴露的核心接口是异步上下文管理器get_session()session_manager.py它把所有并发控制收敛在一个asyncio.Lock内async def __init__(self, max_concurrent_sessions: int 10): self.max_concurrent_sessions max_concurrent_sessions self._sessions: Dict[str, SessionInfo] {} self._computer_pool ComputerPool() self._session_lock asyncio.Lock()关键行为逐条说明会话 ID 缺省生成未传session_id时用str(uuid.uuid4())兜底保证旧调用方零改动可用会话复用session_id已存在时复用同一 Computer 实例并刷新last_activity若该会话处于is_shutting_down状态则抛出RuntimeError避免“清理中会话”被再次使用资源上限会话数达到max_concurrent_sessions默认 10时抛出Maximum concurrent sessions (10) reached从源头防止 VM 数量失控任务注册/注销register_task()/unregister_task()维护active_tasks集合是“有活跃任务就不清理”这一安全语义的数据基础主动清理cleanup_session()发现active_tasks非空时只标记is_shutting_down True而返回等任务跑完否则调用_force_cleanup_session()把 Computer 归还池中并删除会话后台空闲回收_cleanup_loop()每 60 秒扫描一次将“无活跃任务且空闲超过 600 秒10 分钟”的会话强制清理session_manager.py。注意一个实现细节cleanup_idle()池层面目前是空实现占位源码注释说明“well keep instances in pool”——真正生效的自动清理发生在会话层的_cleanup_loop。阅读源码时以这里为准。3. 服务器工具层session_id 参数与向后兼容server.py 中所有工具都注册在FastMCP(namecua-agent)上并统一支持可选的session_id参数签名与文档一致server.tool(structured_outputFalse) async def screenshot_cua(ctx: Context, session_id: Optional[str] None) - Any: ... server.tool(structured_outputFalse) async def run_cua_task(ctx: Context, task: str, session_id: Optional[str] None) - Any: ... server.tool(structured_outputFalse) async def run_multi_cua_tasks( ctx: Context, tasks: List[str], session_id: Optional[str] None, concurrent: bool False ) - Any: ...此外还有两个文档“Session Management”示例用到的管理工具server.pyget_session_stats(ctx)返回total_sessions、max_concurrent和每个会话的created_at、last_activity、active_tasks、is_shutting_down统计结构见 session_manager.pycleanup_session(ctx, session_id)发起指定会话的清理返回Session {session_id} cleanup initiated。3.1 run_cua_task 内部的会话协作以run_cua_task为例可以看到会话、任务注册与错误处理的完整协作链server.py生成task_id str(uuid.uuid4())进入async with session_manager.get_session(session_id) as sessionawait session_manager.register_task(session.session_id, task_id)把任务计入会话用ComputerAgent(modelmodel_name, only_n_most_recent_images..., tools[session.computer])创建 agentagent 绑定的工具是该会话自己的 Computer 实例——这正是资源隔离的落点async for result in agent.run(messages)逐条流式处理message/tool_use/tool_result输出通过ctx.yield_message/yield_tool_call/yield_tool_output上报给 MCP 客户端——文档“Streaming updates prevent timeout issues”提到的超时问题即由此缓解finally块中unregister_task确保异常路径也不会泄漏任务记录若整体抛出异常还会尝试对同一session_id再取一次会话、截取一张错误现场截图返回取不到时返回空数据的占位Image。3.2 并发任务执行run_multi_cua_tasks提供两种模式server.py顺序模式默认逐个await run_cua_task(ctx, task, session_id)每完成一个任务调用ctx.report_progress((i 1) / total_tasks)并发模式concurrentTrue为每个任务构建带进度上报的协程然后await asyncio.gather(*task_coroutines, return_exceptionsTrue)。return_exceptionsTrue是关键单个任务失败会以Exception对象返回代码把它替换为(fTask failed: {str(result)}, Image(formatpng, datab))占位结果从而“单个任务失败不阻断其他任务”且结果顺序与输入任务顺序保持一致。值得强调的一点并发模式下的多个任务传入的是同一个session_id它们共享同一会话的 Computer 实例并行操作——这与“多客户端隔离”是两层不同的语义使用concurrentTrue时应注意多个 agent 在同一台 VM 上并行操作可能互相影响。4. 优雅关闭信号处理与资源回收文档第 5 节描述的 Graceful Shutdown 在源码中对应三处server.py信号处理run_server()内注册signal.signal(signal.SIGINT, signal_handler)与SIGTERM收到信号后asyncio.create_task(graceful_shutdown())关闭链graceful_shutdown()调用shutdown_session_manager()→SessionManager.stop()后者先取消清理协程再_force_cleanup_session逐个清理所有会话最后ComputerPool.shutdown()关闭全部实例session_manager.py兜底清理run_server()的finally块保证即使启动阶段异常也会执行await shutdown_session_manager()进程入口用anyio.run(run_server)代替asyncio.run注释说明目的是“avoid nested event loop issues”。使用层面即文档所述# Send SIGTERM for graceful shutdown kill -TERM server_pid # Or use CtrlC (SIGINT)5. 使用示例以下示例完整继承自 CONCURRENT_SESSIONS.md均与 server.py 中的实际工具签名一致。5.1 基础用法向后兼容# These calls work exactly as before await screenshot_cua(ctx) await run_cua_task(ctx, Open browser) await run_multi_cua_tasks(ctx, [Task 1, Task 2])5.2 多客户端隔离用法# Client 1 session_id_1 client-1-session await screenshot_cua(ctx, session_id_1) await run_cua_task(ctx, Open browser, session_id_1) # Client 2 (completely isolated) session_id_2 client-2-session await screenshot_cua(ctx, session_id_2) await run_cua_task(ctx, Open editor, session_id_2)两个客户端各自由SessionManager分配或从池中复用独立的 Computer 实例互不干扰。5.3 并发任务执行# Run tasks concurrently instead of sequentially tasks [Open browser, Open editor, Open terminal] results await run_multi_cua_tasks(ctx, tasks, concurrentTrue)5.4 会话管理# Get session statistics stats await get_session_stats(ctx) print(fActive sessions: {stats[total_sessions]}) # Cleanup specific session await cleanup_session(ctx, session-to-cleanup)6. 配置说明6.1 环境变量环境变量作用默认值源码位置CUA_MODEL_NAMEAgent 使用的模型anthropic/claude-sonnet-4-5-20250929server.pyCUA_MAX_IMAGESagent 保留的最大历史截图数3以 int 解析server.pyCUA_USE_HOST_COMPUTER_SERVER池中实例是否连接宿主机 computer-servertrue/1/yes生效falsesession_manager.py6.2 会话管理器与实例池参数# In session_manager.py class SessionManager: def __init__(self, max_concurrent_sessions: int 10): # 最大并发会话数 class ComputerPool: def __init__(self, max_size: int 5, idle_timeout: float 300.0): # 最大实例数与空闲超时源码中还存在两个未被环境变量暴露的硬编码常量会话空闲回收超时idle_timeout 600.010 分钟与清理扫描间隔 60 秒session_manager.py。当前版本如需调整只能修改源码这与文档“Future Enhancements”中提到的“可配置会话超时”方向一致。运行环境方面pyproject.toml 声明requires-python 3.12,3.14依赖mcp1.6.0,2.0.0、cua-agent[all]0.8.0、cua-computer0.4.0,0.5.0入口命令为cua-mcp-server映射到mcp_server.server:main。7. 测试验证文档声明的测试文件在仓库根目录的 tests/test_mcp_server_session_management.py。该测试通过 stub 掉mcp.server.fastmcp、computer、cua_agent等外部依赖后用importlib动态加载真实的 server.py 执行断言覆盖了文档“Testing”一节列出的全部能力会话创建与复用test_screenshot_cua_creates_new_session、test_session_reuse_with_same_id并发会话隔离test_concurrent_sessions_isolation两个不同session_id的任务用asyncio.gather并行跑顺序与并发多任务test_run_multi_cua_tasks_sequential、test_run_multi_cua_tasks_concurrent统计与清理test_get_session_stats、test_cleanup_session错误处理test_error_handling_with_session_management模拟 agent 抛RuntimeError断言返回Error during task execution前缀且仍带 PNG 结果。运行方式与文档一致从仓库根目录执行pytest tests/test_mcp_server_session_management.py -v另有两个相关测试可作补充参考tests/test_mcp_server_streaming.py 验证流式输出路径libs/python/mcp-server/tests/test_mcp_server.py 做包级导入冒烟测试。8. 迁移指南8.1 存量客户端零改动# This still works exactly as before await run_cua_task(ctx, My task)不传session_id时自动使用随机 UUID 会话语义上与旧的单客户端行为兼容隔离粒度变为每次调用一个会话资源由池管理。8.2 新的多客户端应用# Create a unique session ID for each client session_id str(uuid.uuid4()) await run_cua_task(ctx, My task, session_id)8.3 并发任务tasks [Task 1, Task 2, Task 3] results await run_multi_cua_tasks(ctx, tasks, concurrentTrue)9. 监控与日志会话统计get_session_stats()返回total_sessions、max_concurrent以及每个会话的created_at/last_activity/active_tasks/is_shutting_down可周期性轮询观察池与负载水位日志server.py 把日志级别设为DEBUG并输出到 stderr覆盖会话创建/清理、任务注册/完成、池使用量与错误恢复等关键事件logger 名分别为mcp-server与mcp-server.session_manager关闭kill -TERM pid或 CtrlC 触发第 4 节所述的完整清理链。10. 改造收益与后续方向对照文档的性能对比表改造前后差异为维度改造前改造后Computer 实例单一全局实例每会话独立实例 池化复用多客户端相互干扰、资源冲突会话级隔离任务执行仅顺序concurrentTrue并行执行关闭无清理信号驱动的优雅关闭长任务30s 超时问题流式更新规避超时资源上限无可配置会话/池上限 自动空闲回收文档同时列出后续增强方向会话持久化、跨实例负载均衡、实时资源监控、池容量自动伸缩、按会话类型可配置的超时。结合源码现状cleanup_idle()空实现、空闲超时硬编码这些方向都有明确的落点。小结Cua MCP Server 通过SessionManager会话编排、上限控制、自动回收ComputerPool实例复用、生命周期管理 统一session_id参数三个层次把一个单客户端的 MCP 服务改造为可多客户端并发、可优雅关闭的服务。所有关键行为均可在 session_manager.py、server.py 中逐一对照源码验证并通过 tests/test_mcp_server_session_management.py 中的测试用例复现。【免费下载链接】cuaScale computer-use 2.0 with open-source drivers, cross-OS fleets, and benchmarks for training, evaluation, and data generation.项目地址: https://gitcode.com/GitHub_Trending/cua/cua创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表