OpenManus源码分析-Flow
Open Manus源码架构分析
Open Manus源码分析-Agent
在前两篇文章中,我们分析了OpenManus的整体架构和核心Agent模块。今天我们深入探讨OpenManus的Flow模块,这是整个框架中负责编排工作流的关键组件。通过本文,你将了解OpenManus如何组织复杂的任务流程,以及如何利用其灵活的流程设计构建智能应用。
概述
Open Manus Flow 是对Agent模型的封装,用作编排多个Agent的调度,完成用户指令。简单来说,Flow就像是一个"导演",它协调不同的Agent,按照预设的逻辑顺序完成复杂任务。
编排Flow的提示词为:
System Prompt:
"You are a planning assistant. Create a concise, actionable plan with clear steps. "
"Focus on key milestones rather than detailed sub-steps. "
“Optimize for clarity and efficiency.”
User Prompt:
“Create a reasonable plan with clear steps to accomplish the task: {request}”
request为用户输入提问。
Tool Planning:
parameters: dict = {
“type”: “object”,
“properties”: {
“command”: {
“description”: “The command to execute. Available commands: create, update, list, get, set_active, mark_step, delete.”,
“enum”: [
“create”,
“update”,
“list”,
“get”,
“set_active”,
“mark_step”,
“delete”,
],
“type”: “string”,
},
“plan_id”: {
“description”: “Unique identifier for the plan. Required for create, update, set_active, and delete commands. Optional for get and mark_step (uses active plan if not specified).”,
“type”: “string”,
},
“title”: {
“description”: “Title for the plan. Required for create command, optional for update command.”,
“type”: “string”,
},
“steps”: {
“description”: “List of plan steps. Required for create command, optional for update command.”,
“type”: “array”,
“items”: {“type”: “string”},
},
“step_index”: {
“description”: “Index of the step to update (0-based). Required for mark_step command.”,
“type”: “integer”,
},
“step_status”: {
“description”: “Status to set for a step. Used with mark_step command.”,
“enum”: [“not_started”, “in_progress”, “completed”, “blocked”],
“type”: “string”,
},
“step_notes”: {
“description”: “Additional notes for a step. Optional for mark_step command.”,
“type”: “string”,
},
},
“required”: [“command”],
“additionalProperties”: False,
}
核心组件解析

通过源码分析,我们可以看到Flow模块主要包含以下核心组件:
1. BaseFlow
BaseFlow是所有流程的基类,定义了流程的基本结构和行为。让我们看看它的核心定义:
class BaseFlow(BaseModel, ABC):
"""支持多个代理的执行流程基类"""
agents: Dict[str, BaseAgent] # 代理字典,键为代理名称,值为代理实例
tools: Optional[List] = None # 可选的工具列表
primary_agent_key: Optional[str] = None # 主代理的键名
def __init__(
self, agents: Union[BaseAgent, List[BaseAgent], Dict[str, BaseAgent]], **data
):
# 处理不同方式提供的代理
if isinstance(agents, BaseAgent):
# 如果只提供了单个代理,将其作为"default"代理
agents_dict = {"default": agents}
elif isinstance(agents, list):
# 如果提供了代理列表,使用索引作为键名
agents_dict = {f"agent_{i}": agent for i, agent in enumerate(agents)}
else:
# 如果已经是字典,直接使用
agents_dict = agents
# 如果未指定主代理,使用第一个代理作为主代理
primary_key = data.get("primary_agent_key")
if not primary_key and agents_dict:
primary_key = next(iter(agents_dict))
data["primary_agent_key"] = primary_key
# 设置代理字典
data["agents"] = agents_dict
# 使用BaseModel的初始化方法
super().__init__(**data)
@property
def primary_agent(self) -> Optional[BaseAgent]:
"""获取流程的主代理"""
return self.agents.get(self.primary_agent_key)
@abstractmethod
async def execute(self, input_text: str) -> str:
"""执行流程,处理给定的输入文本
必须由子类实现具体执行逻辑"""
这个基类定义了流程的基本框架,包括状态管理、步骤控制和执行循环。
2. PlanningFlow
PlanningFlow是目前OpenManus中实现的主要流程类型,它专注于任务规划和执行:
class PlanningFlow(BaseFlow):
"""使用代理管理任务规划和执行的流程"""
llm: LLM = Field(default_factory=lambda: LLM()) # 用于任务规划的语言模型
planning_tool: PlanningTool = Field(default_factory=PlanningTool) # 计划工具
executor_keys: List[str] = Field(default_factory=list) # 执行代理的键名列表
active_plan_id: str = Field(default_factory=lambda: f"plan_{int(time.time())}") # 当前计划ID
current_step_index: Optional[int] = None # 当前执行步骤的索引
async def execute(self, input_text: str) -> str:
"""执行规划流程,协调代理完成任务"""
try:
if not self.primary_agent:
raise ValueError("没有可用的主代理")
# 如果提供了输入文本,创建初始计划
if input_text:
await self._create_initial_plan(input_text)
result = ""
while True:
# 获取当前要执行的步骤
self.current_step_index, step_info = await self._get_current_step_info()
# 如果没有更多步骤或计划已完成,退出循环
if self.current_step_index is None:
result += await self._finalize_plan()
break
# 根据步骤类型选择合适的执行代理
step_type = step_info.get("type") if step_info else None
executor = self.get_executor(step_type)
# 执行步骤并获取结果
step_result = await self._execute_step(executor, step_info)
result += step_result + "\n"
# 检查代理是否希望终止执行
if hasattr(executor, "state") and executor.state == AgentState.FINISHED:
break
return result
except Exception as e:
logger.error(f"PlanningFlow执行错误: {str(e)}")
return f"执行失败: {str(e)}"
3. FlowFactory
FlowFactory提供了创建不同类型Flow的工厂方法,便于应用层使用:
class FlowType(str, Enum):
PLANNING = "planning" # 规划类型流程
class FlowFactory:
"""创建不同类型流程的工厂类,支持多代理"""
@staticmethod
def create_flow(
flow_type: FlowType,
agents: Union[BaseAgent, List[BaseAgent], Dict[str, BaseAgent]],
**kwargs,
) -> BaseFlow:
# 流程类型映射表
flows = {
FlowType.PLANNING: PlanningFlow,
}
# 获取对应的流程类
flow_class = flows.get(flow_type)
if not flow_class:
raise ValueError(f"未知的流程类型: {flow_type}")
# 创建并返回流程实例
return flow_class(agents, **kwargs)
PlanningFlow深度解析
PlanningFlow是OpenManus当前实现的主要Flow类型,它实现了一个基于规划的工作流程。让我们深入了解它的工作机制:
1. 任务规划
PlanningFlow首先会通过LLM创建一个初始计划:
async def _create_initial_plan(self, request: str) -> None:
"""使用流程的LLM和PlanningTool根据请求创建初始计划"""
logger.info(f"创建初始计划,ID: {self.active_plan_id}")
# 创建计划助手的系统消息
system_message = Message.system_message(
"你是一个规划助手。创建一个简洁、可操作的计划,包含明确的步骤。"
"专注于关键里程碑而非详细的子步骤。"
"优化清晰度和效率。"
)
# 创建包含请求的用户消息
user_message = Message.user_message(
f"创建一个合理的计划,包含明确的步骤,以完成以下任务: {request}"
)
# 调用LLM使用PlanningTool
response = await self.llm.ask_tool(
messages=[user_message],
system_msgs=[system_message],
tools=[self.planning_tool.to_param()],
tool_choice=ToolChoice.AUTO,
)
# 处理工具调用结果
if response.tool_calls:
for tool_call in response.tool_calls:
if tool_call.function.name == "planning":
# 解析参数并执行工具
args = json.loads(tool_call.function.arguments)
args["plan_id"] = self.active_plan_id
result = await self.planning_tool.execute(**args)
return
2. 步骤执行
计划创建后,PlanningFlow会逐步执行计划中的每个步骤:
async def _execute_step(self, executor: BaseAgent, step_info: dict) -> str:
"""使用给定的代理执行计划中的单个步骤"""
step_text = step_info.get("text", "")
logger.info(f"执行步骤: {step_text}")
# 创建包含计划信息的上下文
plan_text = await self._get_plan_text()
context = f"当前计划:\n{plan_text}\n\n当前要执行的步骤: {step_text}"
# 使用代理执行步骤
result = await executor.run(context)
# 标记步骤为已完成
await self._mark_step_completed()
return f"步骤结果: {result}"
3. 计划状态管理
PlanningFlow使用PlanStepStatus枚举来管理步骤状态:
class PlanStepStatus(str, Enum):
"""定义计划步骤可能状态的枚举类"""
NOT_STARTED = "not_started" # 未开始
IN_PROGRESS = "in_progress" # 进行中
COMPLETED = "completed" # 已完成
BLOCKED = "blocked" # 被阻塞
@classmethod
def get_all_statuses(cls) -> list[str]:
"""返回所有可能的步骤状态值列表"""
return [status.value for status in cls]
@classmethod
def get_active_statuses(cls) -> list[str]:
"""返回表示活动状态的值列表(未开始或进行中)"""
return [cls.NOT_STARTED.value, cls.IN_PROGRESS.value]
@classmethod
def get_status_marks(cls) -> Dict[str, str]:
"""返回状态到标记符号的映射"""
return {
cls.COMPLETED.value: "[✓]", # 已完成标记
cls.IN_PROGRESS.value: "[→]", # 进行中标记
cls.BLOCKED.value: "[!]", # 阻塞标记
cls.NOT_STARTED.value: "[ ]", # 未开始标记
}
这允许Flow跟踪每个步骤的执行状态,决定下一步操作。
PlanningTool详解
PlanningFlow的核心是与PlanningTool的集成,这个工具提供了计划的创建和管理功能:
class PlanningTool(BaseTool):
"""用于创建和管理执行计划的工具"""
name: str = "planning" # 工具名称
description: str = "创建和管理执行计划" # 工具描述
plans: Dict[str, Dict] = Field(default_factory=dict) # 存储计划的字典
parameters: dict = {
"type": "object",
"required": ["command"],
"properties": {
"command": {
"type": "string",
"enum": ["create", "update", "mark_step", "get_plan"],
"description": "要执行的规划命令",
},
"plan_id": {
"type": "string",
"description": "计划的唯一标识符",
},
"title": {
"type": "string",
"description": "计划标题",
},
"steps": {
"type": "array",
"items": {"type": "string"},
"description": "计划中的步骤",
},
"step_index": {
"type": "integer",
"description": "要标记的步骤索引",
},
"step_status": {
"type": "string",
"enum": ["not_started", "in_progress", "completed", "blocked"],
"description": "设置步骤的状态",
},
},
}
async def execute(self, **kwargs) -> ToolResult:
"""执行规划工具命令"""
command = kwargs.get("command")
if command == "create":
# 创建新计划
return await self._create_plan(**kwargs)
elif command == "update":
# 更新现有计划
return await self._update_plan(**kwargs)
elif command == "mark_step":
# 标记步骤状态
return await self._mark_step(**kwargs)
elif command == "get_plan":
# 获取计划信息
return await self._get_plan(**kwargs)
else:
# 未知命令
return ToolFailure(error=f"未知的规划命令: {command}")
工作流实例分析
为了更直观地理解Flow模块的工作方式,让我们通过一个实际例子来分析。以下是run_flow.py中的使用示例:
async def run_flow():
# 创建代理字典,这里只有一个Manus代理
agents = {
"manus": Manus(),
}
# 获取用户输入
prompt = input("输入你的提示: ")
# 使用工厂方法创建规划流程
flow = FlowFactory.create_flow(
flow_type=FlowType.PLANNING,
agents=agents,
)
logger.warning("正在处理你的请求...")
# 记录开始时间
start_time = time.time()
# 执行流程
result = await flow.execute(prompt)
# 计算执行时间
elapsed_time = time.time() - start_time
logger.info(f"请求处理耗时 {elapsed_time:.2f} 秒")
logger.info(result)
在这个例子中:
- 创建了一个Manus代理
- 通过FlowFactory创建了一个PlanningFlow
- 执行flow.execute(),传入用户提示
- Flow内部创建计划并执行每个步骤
- 返回最终结果
Flow执行

执行过程大致如下:
初始化:用户输入任务描述
创建计划:通过LLM和PlanningTool创建初始任务计划
执行循环:
- 获取当前要执行的步骤
- 根据步骤类型选择合适的执行Agent
- 执行步骤并获取结果
- 标记步骤为已完成
- 检查是否还有下一步
完成:处理最终结果并返回
总结
OpenManus的Flow模块提供了一个强大而灵活的框架,用于组织和协调复杂的智能工作流程。尽管目前主要实现了PlanningFlow,但其基础架构设计允许未来扩展更多类型的Flow。
Flow模块的核心优势包括:
- 灵活的Agent管理:支持多种方式定义和使用多个Agent
- 任务规划能力:通过LLM创建和管理动态计划
- 状态追踪:完善的步骤状态管理
- 可扩展性:基于BaseFlow可以实现各种定制流程
- 工厂模式:通过FlowFactory简化Flow创建
尽管与最初设想的多种Flow实现相比,当前的OpenManus Flow模块相对简单,但它已经提供了处理复杂任务的基本框架。通过BaseFlow的抽象设计,未来可以轻松扩展更多专用Flow类型,如条件分支Flow或并行执行Flow。
通过Flow模块,OpenManus实现了复杂任务的智能编排,为构建下一代AI应用提供了强大的基础设施。
更多推荐


所有评论(0)