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/cua
MCP(Model Context Protocol)Server 让 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.py):
@dataclass 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_factory=set) is_shutting_down: bool = False字段语义直接对应文档中的会话生命周期四阶段(创建 → 任务注册 → 活动追踪 → 清理):last_activity用于空闲判定,active_tasks防止清理有活跃任务的会话,is_shutting_down则拦截对正在清理的会话的新请求。
2.2 ComputerPool:Computer 实例池
ComputerPool负责 Computer 实例的复用与回收,默认参数为max_size: int = 5, idle_timeout: float = 300.0(session_manager.py):
class 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(verbosity=logging.INFO, use_host_computer_server=use_host)即通过环境变量CUA_USE_HOST_COMPUTER_SERVER可让池中的实例连接到宿主机上已有的 computer-server,而不是由池自行启动,这是一个文档未提及但从源码可直接确认的部署选项; 3.轮询等待:池满时进入while not self._available: await asyncio.sleep(0.1)的等待循环,直到有实例被释放回来。
release()把实例从_in_use移回_available;shutdown()则对两个集合中的实例统一调用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()(池层面)目前是空实现占位,源码注释说明“we'll keep instances in pool”——真正生效的自动清理发生在会话层的_cleanup_loop。阅读源码时以这里为准。
3. 服务器工具层:session_id 参数与向后兼容
server.py 中所有工具都注册在FastMCP(name="cua-agent")上,并统一支持可选的session_id参数,签名与文档一致:
@server.tool(structured_output=False) async def screenshot_cua(ctx: Context, session_id: Optional[str] = None) -> Any: ... @server.tool(structured_output=False) async def run_cua_task(ctx: Context, task: str, session_id: Optional[str] = None) -> Any: ... @server.tool(structured_output=False) async def run_multi_cua_tasks( ctx: Context, tasks: List[str], session_id: Optional[str] = None, concurrent: bool = False ) -> Any: ...此外还有两个文档“Session Management”示例用到的管理工具(server.py):
get_session_stats(ctx):返回total_sessions、max_concurrent和每个会话的created_at、last_activity、active_tasks、is_shutting_down(统计结构见 session_manager.py);cleanup_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 session; await session_manager.register_task(session.session_id, task_id),把任务计入会话;- 用
ComputerAgent(model=model_name, only_n_most_recent_images=..., tools=[session.computer])创建 agent,agent 绑定的工具是该会话自己的 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); - 并发模式(
concurrent=True):为每个任务构建带进度上报的协程,然后await asyncio.gather(*task_coroutines, return_exceptions=True)。return_exceptions=True是关键:单个任务失败会以Exception对象返回,代码把它替换为(f"Task failed: {str(result)}", Image(format="png", data=b""))占位结果,从而“单个任务失败不阻断其他任务”,且结果顺序与输入任务顺序保持一致。
值得强调的一点:并发模式下的多个任务传入的是同一个session_id,它们共享同一会话的 Computer 实例并行操作——这与“多客户端隔离”是两层不同的语义,使用concurrent=True时应注意多个 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 Ctrl+C (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, concurrent=True)5.4 会话管理
# Get session statistics stats = await get_session_stats(ctx) print(f"Active sessions: {stats['total_sessions']}") # Cleanup specific session await cleanup_session(ctx, "session-to-cleanup")6. 配置说明
6.1 环境变量
| 环境变量 | 作用 | 默认值 | 源码位置 |
|---|---|---|---|
CUA_MODEL_NAME | Agent 使用的模型 | anthropic/claude-sonnet-4-5-20250929 | server.py |
CUA_MAX_IMAGES | agent 保留的最大历史截图数 | 3(以 int 解析) | server.py |
CUA_USE_HOST_COMPUTER_SERVER | 池中实例是否连接宿主机 computer-server(true/1/yes生效) | false | session_manager.py |
6.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.0(10 分钟)与清理扫描间隔 60 秒(session_manager.py)。当前版本如需调整只能修改源码,这与文档“Future Enhancements”中提到的“可配置会话超时”方向一致。
运行环境方面,pyproject.toml 声明requires-python = ">=3.12,<3.14",依赖mcp>=1.6.0,<2.0.0、cua-agent[all]>=0.8.0、cua-computer>=0.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, concurrent=True)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>或 Ctrl+C 触发第 4 节所述的完整清理链。
10. 改造收益与后续方向
对照文档的性能对比表,改造前后差异为:
| 维度 | 改造前 | 改造后 |
|---|---|---|
| Computer 实例 | 单一全局实例 | 每会话独立实例 + 池化复用 |
| 多客户端 | 相互干扰、资源冲突 | 会话级隔离 |
| 任务执行 | 仅顺序 | concurrent=True并行执行 |
| 关闭 | 无清理 | 信号驱动的优雅关闭 |
| 长任务 | 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),仅供参考