☰
openJiuwen 执行器 Runner 深度指南:Agent 与 Workflow 的统一异步执行入口
2026/10/12 4:10:35 网站建设 项目流程
  • 人工智能
  • AI Agent
  • Agent 框架
  • 大模型
  • 工具调用
  • RAG
  • 提示工程
  • 强化学习

【免费下载链接】agent-core

openJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力

项目地址:https://gitcode.com/openJiuwen/agent-core
点击查看免费下载

导读

Runner 是 openJiuwen Agent 开发框架中执行所有核心组件(Workflow、Agent)的统一入口与控制中心,它将复杂的执行逻辑抽象为简洁、一致的异步编程接口。本文以官方文档 执行器Runner 为主体,结合 runner.py、runner_config.py 与单元/系统测试源码,系统讲解 Runner 的单例设计、Agent 与 Workflow 的单次/流式执行、运行配置、资源注册与会话生命周期管理。读完本文,你将能够直接用Runner.run_agent/Runner.run_workflow运行任意内置或自定义的 Agent 与工作流,并理解其背后的调用链与会话管理机制。

一、Runner 是什么:统一入口与控制中心

Runner 将"如何运行一个 Agent、如何运行一个 Workflow"的复杂逻辑集中封装,对外只暴露一套编程接口:

  • Agent 执行:提供标准的异步调用invoke与异步流式调用stream两种执行入口;
  • Workflow 执行:同样提供invoke与stream两种执行入口。

也就是说,无论是 ReActAgent、WorkflowAgent 等内置 Agent,还是用户自定义 Agent、自定义工作流,都可以通过 Runner 这一套 API 统一运行,无需关心会话构造、资源解析等底层细节。

重要说明(单例设计):Runner 是一个单例类,所有方法调用和属性访问都会自动代理到全局的 Runner 实例。无需实例化 Runner,直接通过类名调用即可,例如Runner.start()、Runner.resource_mgr。

二、全局单例与代理机制(源码解析)

单例机制在 runner.py 中实现,由三层构成:

1. 真实实现_RunnerImpl

# openjiuwen/core/runner/runner.py class _RunnerImpl(_TeamRunnerMixin): _DEFAULT_RUNNER_ID = "global" _DEFAULT_AGENT_SESSION_ID = "default_session" _AGENT_CONVERSATION_ID = "conversation_id"

它持有全部运行时状态:资源管理器、本地消息队列、回调框架、anyio 根任务组等。

2. 全局唯一实例GLOBAL_RUNNER

# 模块级直接创建,进程内唯一 GLOBAL_RUNNER = _RunnerImpl(config=DEFAULT_RUNNER_CONFIG)

3. 门面类Runner与_ClassProperty描述符

class _ClassProperty: """Descriptor for class-level properties.""" def __get__(self, obj, objtype=None): return getattr(GLOBAL_RUNNER, self.name) class Runner(_TeamRunnerClassMixin): # Properties resource_mgr: ResourceMgr = _ClassProperty("resource_mgr") pubsub = _ClassProperty("pubsub") dist_pubsub = _ClassProperty("dist_pubsub") callback_framework: AsyncCallbackFramework = _ClassProperty("callback_framework")

Runner自身没有任何实例状态,其类方法与类属性一律转发到GLOBAL_RUNNER。因此在代码中Runner.start()、Runner.run_agent(...)、Runner.resource_mgr与直接操作全局实例完全等价。模块 openjiuwen/core/runner/init.py 还通过__getattr__实现了Runner的懒加载,避免不必要的导入开销。

从源码结构可以推断,这样的设计同时服务于多进程/多 Runner 场景:_RunnerImpl支持传入自定义runner_id与RunnerConfig,而Runner门面则固定指向默认全局实例,两者互不干扰。

三、Runner 公共 API 一览

以下是Runner对外暴露的主要类方法与属性(均基于 runner.py 与 team_runner.py 整理):

API类型说明
Runner.start()async 类方法启动 Runner 及其关联组件(消息队列、checkpointer 等)
Runner.stop()async 类方法停止 Runner 并释放资源
Runner.run_agent(agent, inputs, ...)async 类方法单次执行 Agent(支持实例或注册 ID)
Runner.run_agent_streaming(agent, inputs, ...)async 生成器流式执行 Agent
Runner.run_workflow(workflow, inputs, ...)async 类方法单次执行 Workflow(支持实例或注册 ID)
Runner.run_workflow_streaming(workflow, inputs, ...)async 生成器流式执行 Workflow
Runner.run_agent_team(...)async 类方法执行 Agent 团队(agent_teams / multi_agent 两种路径)
Runner.run_agent_team_streaming(...)async 生成器流式执行 Agent 团队
Runner.spawn_agent(agent_config, inputs, ...)async 类方法在子进程中运行 Agent
Runner.spawn_agent_streaming(...)async 生成器在子进程中流式运行 Agent
Runner.release(session_id, force=False)async 类方法释放指定会话关联的资源(checkpoint、动态表等)
Runner.set_config(config)/Runner.get_config()类方法设置/读取 Runner 配置
Runner.resource_mgr属性资源管理器(workflow/agent/agent_team/tool/model/prompt)
Runner.pubsub属性本地发布订阅消息队列
Runner.dist_pubsub属性分布式消息队列(跨进程通信)
Runner.callback_framework属性异步回调框架

其中run_agent/run_workflow系列是日常开发最常用的入口,下面分别展开。

四、使用 Runner 执行 Agent

Runner 支持所有 Agent 的单次输出执行和流式输出执行,包括 ReActAgent、WorkflowAgent 等内置 Agent,也包括用户自定义的 Agent。下面以一个WorkflowAgent为例,介绍完整执行过程。

4.1 创建 WorkflowAgent 实例

from openjiuwen.core.common.constants.enums import ControllerType from openjiuwen.core.application.workflow_agent import WorkflowAgentConfig, WorkflowAgent from openjiuwen.core.workflow import End, Start, Workflow, WorkflowCard, generate_workflow_key from openjiuwen.core.workflow.workflow_config import WorkflowConfig from openjiuwen.core.runner.runner import Runner def create_agent(runner): # 创建工作流 flow card = WorkflowCard(id="workflow_id", name="简单工作流", version="1", description="this_is_a_demo") flow = Workflow(workflow_config=WorkflowConfig(card=card)) flow.set_start_comp("start", Start(), inputs_schema={"query": "${query}"}) flow.set_end_comp("end", End(), inputs_schema={"result": "${start.query}"}) flow.add_connection("start", "end") # 将 flow 注册到资源管理器(使用 generate_workflow_key 生成正确的 key) # 注意:注册时需要使用 id_version 格式的 key register_card = WorkflowCard( id=generate_workflow_key(card.id, card.version), name=card.name, version=card.version, description=card.description ) runner.resource_mgr.add_workflow(register_card, lambda: flow) # 创建Agent,使用 WorkflowCard 描述工作流输入参数 workflow_card = WorkflowCard( id="workflow_id", version="1", name="简单工作流", description="this_is_a_demo", input_params={"query": {"type": "string"}}, ) workflow_agent_config = WorkflowAgentConfig(id="agent_id", version="1", description="this_is_a_demo", workflows=[workflow_card], controller_type=ControllerType.WorkflowController ) agent = WorkflowAgent(agent_config=workflow_agent_config) return agent # Runner 是单例类,直接使用类名即可 agent = create_agent(Runner)

代码中的关键点说明:

  • WorkflowCard是工作流的元数据卡片(定义见 base.py),包含id、name、version、description,可选input_params描述工作流输入参数,version默认值为空字符串;
  • WorkflowConfig负责承载卡片与执行规格(card、spec、workflow_max_nesting_depth,默认最大嵌套深度 5,取值范围 0~10,见 workflow_config.py);
  • flow.set_start_comp / set_end_comp / add_connection分别注册起始组件、结束组件并连接拓扑,示例中start接收输入query,end将start.query透传为结果;
  • generate_workflow_key(card.id, card.version)生成"{workflow_id}_{workflow_version}"格式的注册 key(源码实现见 base.py)。注册到资源管理器时必须使用该id_version格式,否则 Runner 在按 key 查找工作流时会找不到;
  • WorkflowAgentConfig(定义见 legacy/config.py)的controller_type必须为ControllerType.WorkflowController——workflow_agent.py 中WorkflowAgent.__init__会强制校验,否则抛出NotImplementedError;
  • WorkflowAgent基于ControllerAgent实现,其invoke/stream完全委托给控制器。

4.2 调用 run_agent 直接运行

import asyncio print(asyncio.run(Runner.run_agent(agent=agent, inputs={"conversation_id": "id1", "query": "哈哈"})))

执行结果:

{'output': WorkflowOutput(result={'output': {'result': '哈哈'}}, state= < WorkflowExecutionState.COMPLETED: 'COMPLETED' >), 'result_type': 'answer'}

结果字典包含两个关键字段:

  • output:封装了最终输出WorkflowOutput(包含result与state两个字段,见 base.py 中WorkflowOutput定义)。WorkflowExecutionState有三种取值:COMPLETED(正常完成)、INPUT_REQUIRED(等待用户输入)、ERROR(执行出错);
  • result_type:标识本次执行的结果类型。在系统测试(如 test_hitl_rail_chain_tools.py)中可以看到它可能取"answer"、"interrupt"等值,用于区分正常回答与需要人工确认的中间态。

4.3 会话自动管理:conversation_id 的作用

run_agent的inputs中有一个特殊的conversation_id字段。在 runner.py 的_prepare_agent中:

session_id = inputs.get(self._AGENT_CONVERSATION_ID, session if isinstance(session, str) else self._DEFAULT_AGENT_SESSION_ID)

也就是说,Runner 会优先从inputs["conversation_id"]取会话 ID,其次取显式传入的session字符串,否则回退到默认会话"default_session",并自动调用create_agent_session(...)创建 Agent 会话(_create_agent_session还会从 agent 配置中提取_env注入会话)。这意味着你通常无需手动构造 Session,只要在多次调用时复用同一个conversation_id,即可实现多轮对话的上下文延续。

此外,run_agent对不同类型的 Agent 采用了不同的执行路径(源码 runner.py):

  • RemoteAgent(A2A 远端 Agent):直接invoke(inputs),不创建本地会话;
  • LegacyBaseAgent(含 WorkflowAgent、ControllerAgent 体系):以session=None调用invoke,由控制器自行管理会话生命周期;
  • 其他标准BaseAgent:传入自动创建的agent_session,执行后调用agent_session.post_run()完成会话收尾。

4.4 流式执行 Agent

当需要边生成边输出(如 LLM 流式吐字、逐步汇报执行进度)时,使用run_agent_streaming:

import asyncio from openjiuwen.core.runner import Runner async def main(): agent = create_agent(Runner) async for chunk in Runner.run_agent_streaming( agent=agent, inputs={"conversation_id": "id1", "query": "哈哈"}, ): print(chunk) asyncio.run(main())

run_agent_streaming是异步生成器,逐块产出WorkflowChunk(OutputSchema/CustomSchema/TraceSchema的联合类型),并支持stream_modes参数控制输出类型(详见 5.3 节)。

五、使用 Runner 执行 Workflow

Runner 支持 Workflow 的单次输出执行和流式输出执行。下面通过构建一个简单工作流为例介绍完整过程。

5.1 创建一个 Workflow

from openjiuwen.core.workflow import End, Start, Workflow, WorkflowCard from openjiuwen.core.workflow.workflow_config import WorkflowConfig def build_workflow(name, workflow_id, version): flow = Workflow(workflow_config=WorkflowConfig( card=WorkflowCard(id=workflow_id, name=name, version=version, description="this_is_a_demo"))) flow.set_start_comp("start", Start(), inputs_schema={"query": "${query}"}) flow.set_end_comp("end", End(), inputs_schema={"result": "${start.query}"}) flow.add_connection("start", "end") return flow workflow = build_workflow("test_workflow", "test_workflow", "1")

Workflow的核心构建 API 包括set_start_comp(设置起始组件及输入映射)、set_end_comp(设置结束组件及输出映射)、add_connection(连接两个组件),完整签名见 workflow.py。示例中工作流将输入query从start原样透传到end的result输出。

5.2 调用 run_workflow 直接运行

import asyncio from openjiuwen.core.runner import Runner result = asyncio.run(Runner.run_workflow(workflow=workflow, inputs={"query": "query workflow"})) print(result)

执行结果:

result = {'output': {'result': 'query workflow'}} state = < WorkflowExecutionState.COMPLETED: 'COMPLETED' >

无需再显式构造会话——run_workflow内部会自动创建WorkflowSession。这一行为在_create_workflow_session中实现(runner.py):

  • session为空 → 调用create_workflow_session()新建会话(实现见 session/workflow.py,会话 ID 自动生成 UUID);
  • session为字符串 → 以该字符串为 ID 创建会话;
  • session为 Agent 会话 → 调用AgentSession.create_workflow_session()派生工作流会话;
  • 否则直接复用传入的 WorkflowSession 实例。

run_workflow也支持按注册 ID 运行:_prepare_workflow(runner.py)会将工作流card.id与card.version组合为"{id}_{version}"的 key,通过resource_mgr.get_workflow(workflow_id=key, session=...)取回工作流实例再执行。因此既可以直接传Workflow实例,也可以传注册时的 ID 字符串,例如:

result = await Runner.run_workflow("test_workflow_1", inputs={"query": "query workflow"})

这一用法与单元测试 tests/unit_tests/core/runner/test_runner.py 中的test_run_workflow完全一致:

workflow = self._build_workflow(name, workflow_id, version) Runner.resource_mgr.add_workflow(workflow.card, lambda: workflow) session = create_workflow_session() result = await Runner.run_workflow(workflow_id, inputs={"query": "query workflow"}, session=session) assert result == WorkflowOutput(result={"result": "query workflow"}, state=WorkflowExecutionState.COMPLETED)

从测试断言可以看出,run_workflow的返回对象就是WorkflowOutput(result+state结构),与文档中打印出的{'output': ...}/state展示一致。

5.3 流式执行 Workflow 与流模式

import asyncio from openjiuwen.core.runner import Runner from openjiuwen.core.session.stream import BaseStreamMode async def main(): workflow = build_workflow("test_workflow", "test_workflow", "1") async for chunk in Runner.run_workflow_streaming( workflow=workflow, inputs={"query": "query workflow"}, stream_modes=[BaseStreamMode.OUTPUT, BaseStreamMode.TRACE], ): print(chunk) asyncio.run(main())

stream_modes控制流式产出的数据类型,BaseStreamMode枚举定义在 session/stream/base.py:

模式说明
BaseStreamMode.OUTPUT框架定义的标准流数据(OutputSchema)
BaseStreamMode.TRACE图执行产生的追踪流数据(TraceSchema)
BaseStreamMode.CUSTOM可运行对象自定义的流数据(CustomSchema)

底层上,run_workflow_streaming委托给workflow_instance.stream(...)(workflow.py),它逐块产出WorkflowChunk;而run_workflow内部则收集全部块后,根据是否存在交互块(chunk.type == INTERACTION)判定返回INPUT_REQUIRED还是COMPLETED状态,并在必要时应用工作流执行超时(WORKFLOW_EXECUTE_TIMEOUT环境变量)。

六、Runner 运行配置(RunnerConfig)

Runner 的全局配置通过RunnerConfig管理(定义见 runner_config.py),可通过Runner.set_config(cfg)/Runner.get_config()读写。主要字段:

字段默认值说明
distributed_modeTrue(默认配置中为False)是否启用分布式模式(启用后启动分布式消息队列与回复主题订阅)
distributed_configDistributedConfig()分布式配置(见下)
env_prefix""环境前缀,用于 topic 模板隔离多环境
instance_idUUID 自动生成当前 Runner 实例 ID
checkpointer_configNone检查点配置(如type="redis"时启动时自动初始化 Redis checkpointer)
enable_session_controllerFalse是否启用会话控制器
enable_a2aFalse是否启用 A2A 能力

DistributedConfig的默认值:request_timeout=30.0、max_request_concurrency=10000、agent topic 模板"openjiuwen.single_agent.{agent_id}.{version}"、reply topic 模板"openjiuwen.reply.runner.{instance_id}";MessageQueueConfig支持pulsar与fake两种类型(MessageQueueType枚举),PulsarConfig提供url与max_workers(默认 8)。

值得注意的是,模块级默认配置DEFAULT_RUNNER_CONFIG实际将distributed_mode设为False、消息队列设为fake,即开箱即用为单机模式。当start()时若distributed_mode为真,Runner 会通过MessageQueueFactory创建分布式消息队列、激活ReplyTopicSubscription,并等待本地消息队列启动成功(runner.py)。

七、资源管理器 ResourceMgr:Agent 与 Workflow 的注册中心

Runner.resource_mgr是ResourceMgr实例(源码见 resources_manager/resource_manager.py),管理model、workflow、prompt、tool、agent、agent_team六类资源,提供注册、批量注册、按 ID / tag 查询与删除能力。

对工作流执行而言,最常用的是:

  • add_workflow(card, workflow_provider, tag=None):注册工作流,card.id即为注册 ID;
  • get_workflow(workflow_id, tag=..., session=None):按 ID 或 tag 取回工作流实例;
  • add_agent(card, agent_provider, tag=None, interface_url=None)与get_agent(agent_id, ...):注册与获取 Agent。

关键约束:正如 4.1 节所述,工作流注册 key 必须使用generate_workflow_key(id, version)生成的id_version格式。这是因为_prepare_workflow在按字符串查找时总是拼接"{card.id}_{card.version}"作为 key;如果注册时只用了裸id(例如"workflow_id"),运行时会因 key 不匹配而无法命中。这就是原文档中单独构造register_card并重写id的原因。

八、Runner 生命周期:start / stop / release

虽然示例中直接调用run_agent/run_workflow即可运行,但正式应用中建议遵循完整的生命周期管理:

import asyncio from openjiuwen.core.runner import Runner async def main(): await Runner.start() # 启动 Runner:初始化根任务组、checkpointer、消息队列等 try: # ... 执行 Agent / Workflow / Agent 团队 pass finally: await Runner.stop() # 停止并清理资源 asyncio.run(main())
  • start():启动读写锁管理器、确保根任务组就绪;若配置了checkpointer_config(如 Redis)则初始化默认 checkpointer;若为分布式模式则拉起分布式消息队列与回复订阅。启动失败会回滚已启动的组件;
  • stop():按逆序停止回复订阅与消息队列,释放资源管理器(resource_mgr.release()),停止读写锁与根任务组;
  • release(session_id, force=False):释放指定会话的检查点等资源;对 Agent 团队会话还会自动清理任务表、消息表等动态表;force=True时强制停止该会话上仍在活动的团队,否则会抛出AGENT_TEAM_BUSY_INVALID异常。

系统测试 test_multi_workflow_agent.py 中即采用了asyncSetUp中await Runner.start()、asyncTearDown中await Runner.stop()的标准用法。

九、进阶能力:Agent 团队执行与子进程执行

Runner 不止于单 Agent / 单工作流,还作为团队运行与子进程运行的控制中心:

1. Agent 团队执行(team_runner.py):

# 默认路径:agent_teams 的 TeamAgent(传 TeamAgentSpec 或 team_name 字符串) await Runner.run_agent_team(agent_team=spec, inputs={...}) # multi_agent 的 BaseTeam 路径(传 BaseTeam 实例或 team_id) await Runner.run_agent_team(agent_team=team, inputs={...}, base=True) # 已构建的团队成员实例路径(跳过池与激活,仅供 spawn 场景) await Runner.run_agent_team(agent_team=member_agent, inputs={...}, member=True)

配套能力还包括interact_agent_team(向活动团队投递交互消息)、pause_agent_team/stop_agent_team、get_agent_team_monitor、list_active_teams等。

2. 子进程执行:Runner.spawn_agent/Runner.spawn_agent_streaming接收SpawnAgentConfig,在独立子进程中运行 Agent,返回SpawnedProcessHandle用于管理进程(配合spawn_config可启动健康检查),流式版本逐块产出(handle, message)。

这些能力使 Runner 成为横跨单 Agent、工作流、团队、多进程的"统一执行控制面"。

十、测试验证与深入学习

仓库中已有大量针对 Runner 的测试可以对照学习:

  • 单元测试tests/unit_tests/core/runner/test_runner.py:覆盖run_workflow按 ID 执行、WorkflowOutput返回结构断言,以及@tool装饰的工具调用;
  • 系统测试tests/system_tests/agent/workflow_agent/test_multi_workflow_agent.py:多工作流意图路由、打断/恢复场景,展示Runner.start/stop的标准生命周期用法;
  • 系统测试tests/system_tests/agent/react_agent/interrupt/test_hitl_rail_chain_tools.py:展示Runner.run_agent的多轮调用、result_type == "interrupt"的人工确认循环,以及conversation_id保持会话连续性的实际用法。

建议阅读顺序:官方文档 执行器Runner → runner.py(单例与执行核心)→ runner_config.py(配置项)→ resource_manager.py(资源注册)→ 对应测试文件(行为验证)。

结语

Runner 以"单例门面 + 全局实例 + 统一异步接口"的设计,将 Agent、Workflow 乃至 Agent 团队的执行入口收敛到同一套 API:单次执行用run_agent/run_workflow,流式输出用run_agent_streaming/run_workflow_streaming,会话由conversation_id与内部 Session 工厂自动管理,资源通过resource_mgr按id_version规则注册与解析。理解 Runner 的调用链,就等于掌握了 openJiuwen 应用从"定义组件"到"运行组件"之间最关键的一环。

  • 人工智能
  • AI Agent
  • Agent 框架
  • 大模型
  • 工具调用
  • RAG
  • 提示工程
  • 强化学习

【免费下载链接】agent-core

openJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力

项目地址:https://gitcode.com/openJiuwen/agent-core
点击查看免费下载

相关推荐

上一篇:WrenAI完全指南:如何为AI智能体构建数据上下文层的终极解决方案
下一篇:VisualGGPK2终极指南:10分钟掌握《流放之路》资源编辑神器

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询