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)

在这个例子中:

  1. 创建了一个Manus代理
  2. 通过FlowFactory创建了一个PlanningFlow
  3. 执行flow.execute(),传入用户提示
  4. Flow内部创建计划并执行每个步骤
  5. 返回最终结果

Flow执行

在这里插入图片描述
执行过程大致如下:
初始化:用户输入任务描述
创建计划:通过LLM和PlanningTool创建初始任务计划
执行循环:

  1. 获取当前要执行的步骤
  2. 根据步骤类型选择合适的执行Agent
  3. 执行步骤并获取结果
  4. 标记步骤为已完成
  5. 检查是否还有下一步

完成:处理最终结果并返回

总结

OpenManus的Flow模块提供了一个强大而灵活的框架,用于组织和协调复杂的智能工作流程。尽管目前主要实现了PlanningFlow,但其基础架构设计允许未来扩展更多类型的Flow。
Flow模块的核心优势包括:

  1. 灵活的Agent管理:支持多种方式定义和使用多个Agent
  2. 任务规划能力:通过LLM创建和管理动态计划
  3. 状态追踪:完善的步骤状态管理
  4. 可扩展性:基于BaseFlow可以实现各种定制流程
  5. 工厂模式:通过FlowFactory简化Flow创建

尽管与最初设想的多种Flow实现相比,当前的OpenManus Flow模块相对简单,但它已经提供了处理复杂任务的基本框架。通过BaseFlow的抽象设计,未来可以轻松扩展更多专用Flow类型,如条件分支Flow或并行执行Flow。
通过Flow模块,OpenManus实现了复杂任务的智能编排,为构建下一代AI应用提供了强大的基础设施。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐