将 AI 代理添加到工作流时,需要将其包装在执行器中,以便工作流引擎可以将消息路由到工作流、管理其会话状态以及处理其输出。 代理执行程序是处理此适应的内置执行程序。
概述
代理执行程序弥合了代理抽象与工作流执行模型之间的差距。 该方法:
- 从工作流图接收类型化消息,并将其转发到基础代理。
- 管理代理在不同运行之间的会话和对话状态。
- 根据工作流执行模式(流式处理或非流式处理)调整其行为。
- 向工作流调用方生成输出事件(
AgentResponse或AgentResponseUpdate),供其观察。 - 将消息发送到连接的下游执行程序,以便在图形中继续处理。
- 支持为长时间运行的工作流创建检查点。
工作原理
在 C# 中,工作流引擎会在内部为添加到工作流的每个AIAgentHostExecutor创建一个AIAgent。 此专用执行程序扩展 ChatProtocolExecutor 并使用 轮次令牌 模式:
-
消息缓存 - 当消息从其他执行程序到达时,代理执行程序会收集它们。 如果
ForwardIncomingMessages已启用(默认值),则传入消息也会转发到下游执行程序。 -
轮次令牌触发器 — 代理仅在收到
TurnToken令牌后处理其缓存的消息。 -
代理调用 — 执行程序对基础代理调用
RunAsync(非流式处理)或RunStreamingAsync(流式处理)。 -
输出生成 - 如果启用了流式处理事件,则会生成每个增量
AgentResponseUpdate作为工作流输出。 如果EmitAgentResponseEvents已启用,则聚合AgentResponse也会作为工作流输出生成。 - 下游消息传送 - 代理的响应消息将发送到连接的下游执行程序。
-
轮次令牌传递 - 完成轮次后,执行器发送新的
TurnToken到下游,以便链中的下一个代理可以开始运算。
小窍门
某些方案可能需要更专用的代理执行程序;例如, 交接业务流程 使用专用 HandoffAgentExecutor 的自定义路由逻辑。
隐式创建与显式创建
传递AIAgent给WorkflowBuilder时,框架会自动将其包装在AIAgentBinding中,从而创建底层的AIAgentHostExecutor。 无需直接实例化代理执行程序。
AIAgent writerAgent = /* create your agent */;
AIAgent reviewerAgent = /* create your agent */;
// Agents are automatically wrapped — no manual executor creation required
var workflow = new WorkflowBuilder(writerAgent)
.AddEdge(writerAgent, reviewerAgent)
.Build();
可以使用 AgentWorkflowBuilder 上的辅助方法来实现常见模式:
// Build a sequential pipeline of agents
var workflow = AgentWorkflowBuilder.BuildSequential(writerAgent, reviewerAgent);
自定义配置
若要自定义代理执行程序的行为方式,请使用 BindAsExecutorAIAgentHostOptions:
var options = new AIAgentHostOptions
{
EmitAgentUpdateEvents = true,
EmitAgentResponseEvents = true,
ReassignOtherAgentsAsUsers = true,
ForwardIncomingMessages = true,
};
ExecutorBinding writerBinding = writerAgent.BindAsExecutor(options);
var workflow = new WorkflowBuilder(writerBinding)
.AddEdge(writerBinding, reviewerAgent)
.Build();
输入类型
C# 中的代理执行程序接受多个输入类型: string、 ChatMessage和 IEnumerable<ChatMessage>。 字符串输入会自动转换为具有 ChatMessage 角色的 User 实例。 所有传入的消息都会累积,直到收到TurnToken为止,此时执行器会处理消息批次。 启用时 ReassignOtherAgentsAsUsers (默认),其他代理的消息将重新分配给 User 角色,以便基础模型将它们视为用户输入,而来自当前代理的消息将保留该 Assistant 角色。
输出和链式调用
代理完成轮次后,执行者:
- 将代理的响应消息发送到所有连接的下游执行程序。
- 转发新的
TurnToken,以便链中的下一个代理可以开始处理。
这使得链接代理非常简单 , 只需使用边缘连接它们:
var workflow = new WorkflowBuilder(frenchTranslator)
.AddEdge(frenchTranslator, spanishTranslator)
.AddEdge(spanishTranslator, englishTranslator)
.Build();
流式处理行为
流式处理行为由 EmitAgentUpdateEvents 选项在 AIAgentHostOptions 上控制,或通过 TurnToken 动态控制。
-
启用后 , 执行程序对代理调用
RunStreamingAsync并生成每个AgentResponseUpdate代理作为工作流输出事件。 这提供逐个令牌的实时更新。 -
禁用时 , 执行程序调用
RunAsync并生成单个完整响应。
// Enable streaming events at the configuration level
var options = new AIAgentHostOptions
{
EmitAgentUpdateEvents = true,
};
// Or enable streaming dynamically via TurnToken
await run.TrySendMessageAsync(new TurnToken(emitEvents: true));
共享会话
默认情况下,每个代理执行程序维护自己的会话。 若要在代理之间共享会话,请在将代理添加到工作流之前使用通用会话提供程序配置代理。
配置选项
AIAgentHostOptions 控制代理执行程序的行为:
| 选项 | 默认 | 说明 |
|---|---|---|
EmitAgentUpdateEvents |
null |
在执行期间发出流式更新事件。
TurnToken 如果已设置,则优先。 如果两者都是null,则禁用流式传输。 |
EmitAgentResponseEvents |
false |
以工作流输出事件的形式发出聚合代理响应。 |
InterceptUserInputRequests |
false |
截获 UserInputRequestContent 并将其作为工作流消息进行路由处理。 |
InterceptUnterminatedFunctionCalls |
false |
不带相应结果的截获 FunctionCallContent 并将其路由为工作流消息。 |
ReassignOtherAgentsAsUsers |
true |
将来自其他代理 User 的消息重新分配给角色,以便模型将其视为用户输入。 |
ForwardIncomingMessages |
true |
在代理生成消息之前,将传入消息转发到下游执行器。 |
检查点
代理执行程序支持检查点功能,以便对长期运行的工作流进行处理。 当进行检查点时,执行器会进行序列化:
- 代理的会话状态(通过
SerializeSessionAsync)。 - 当前轮次的事件排放配置(仅在请求挂起且执行程序尚未生成其传入
TurnToken时存在)。 - 任何挂起的用户输入请求和函数调用请求。
还原时,执行器将反序列化会话和挂起的请求状态,从而允许工作流从其离开的位置恢复。
工作原理
AgentExecutor 类封装了一个实现 SupportsAgentRun 协议的代理。 执行者收到消息时:
-
消息规范化 - 输入规范化为对象列表
Message,并添加到执行程序的内部缓存中。 执行程序接受多个输入类型(str、Message、list[str | Message]AgentExecutorRequest和AgentExecutorResponse),每个类型都路由到一个专用处理程序,该处理程序在缓存之前规范化输入。 -
代理调用 — 执行程序使用缓存的消息调用
agent.run(),根据工作流执行模式自动选择流模式或非流模式。 -
输出排放 - 在流式处理模式下,每个
AgentResponseUpdate输出都作为工作流输出事件生成。 在非流式处理模式下,会输出一个AgentResponse。 -
下游调度 - 代理完成后,执行程序会向所有连接的下游执行程序发送一个
AgentExecutorResponse。 此响应包括完整的对话历史记录,可实现无缝链接。 - 缓存重置 - 执行程序的内部消息缓存在调用代理后清除,确保每个代理调用仅处理自上次调用以来收到的新消息。
小窍门
某些方案可能需要更专用的代理执行程序;例如, 交接业务流程 使用具有自定义路由逻辑的专用执行程序。
隐式创建与显式创建
直接传递代理时,WorkflowBuilder 会自动将代理包装在 AgentExecutor 实例中。 对于大多数工作流,隐式创建就足够了:
from agent_framework import WorkflowBuilder
writer_agent = client.as_agent(name="Writer", instructions="...")
reviewer_agent = client.as_agent(name="Reviewer", instructions="...")
# Agents are automatically wrapped — no manual AgentExecutor creation required
workflow = (
WorkflowBuilder(start_executor=writer_agent)
.add_edge(writer_agent, reviewer_agent)
.build()
)
显式创建
在需要以下情况下显式创建:AgentExecutor
- 在多个代理之间共享会话。
- 为路由和目标运行时的 kwargs 提供自定义执行器 ID。
- 引用多个边缘中的同一执行程序实例。
from agent_framework import AgentExecutor
writer_executor = AgentExecutor(writer_agent, id="my-writer")
reviewer_executor = AgentExecutor(reviewer_agent, id="my-reviewer")
workflow = (
WorkflowBuilder(start_executor=writer_executor)
.add_edge(writer_executor, reviewer_executor)
.build()
)
构造函数参数:
| 参数 | 类型 | 说明 |
|---|---|---|
agent |
SupportsAgentRun |
要包装的代理。 |
session |
AgentSession \| None |
用于代理运行的会话。 如果 None,则由代理创建新会话。 |
id |
str \| None |
唯一执行器 ID。 如果可用,则默认为代理的名称。 |
context_mode |
"full" \| "last_agent" \| "custom" \| None |
控制从上游代理接收 AgentExecutorResponse 会话上下文时的处理方式。 默认为 "full",它提供上游代理的完整会话(输入 + 响应)。 请参阅 上下文模式。 |
context_filter |
Callable[[list[Message]], list[Message]] \| None |
用于选择要包含的消息的自定义筛选器函数。 当 context_mode 是 "custom" 时为必需项。 |
小窍门
执行程序 ID 也是当你将 workflow.run(function_invocation_kwargs=...) 或 client_kwargs= 用于定位单个代理时使用的密钥。 如果省略 id,工作流将使用封装代理的名称。
输入类型
定义 AgentExecutor 多个处理程序方法,每个方法都接受不同的输入类型。 工作流引擎根据消息类型自动调度正确的处理程序。 除了AgentExecutorRequest之外,所有输入类型都会触发代理立即运行,其中should_respond标志控制代理是运行还是仅缓存消息。
| 输入类型 | 处理器 | 触发器代理 | 说明 |
|---|---|---|---|
AgentExecutorRequest |
run |
有條件的 | 规范输入类型。 包含消息列表和一个 should_respond 标志,用于控制代理是否运行。 |
str |
from_str |
始终 | 接受一个原始字符串提示。 |
Message |
from_message |
始终 | 接受单个 Message 对象。 |
list[str \| Message] |
from_messages |
始终 | 接受字符串或 Message 对象列表作为对话上下文。 |
AgentExecutorResponse |
from_response |
始终 | 接受先前代理执行器的响应,实现直接链式调用。 |
使用 AgentExecutorRequest
AgentExecutorRequest 是规范输入类型,提供最多的控件:
from agent_framework import AgentExecutorRequest, Message
# Create a request with messages
request = AgentExecutorRequest(
messages=[Message(role="user", contents=["Hello, world!"])],
should_respond=True,
)
# Run the workflow
result = await workflow.run(request)
该 should_respond 标志控制代理是立即处理消息,还是只是缓存消息以供以后使用:
-
True(默认值) - 代理运行并生成响应。 -
False— 消息将添加到缓存,但代理未运行。 这对于在触发响应之前预加载会话上下文非常有用。
输出和链式调用
代理完成后,执行程序将发送 AgentExecutorResponse 到下游。 此数据类包含:
| 领域 | 类型 | 说明 |
|---|---|---|
executor_id |
str |
生成响应的执行程序的 ID。 |
agent_response |
AgentResponse |
原始代理响应(未被客户端更改)。 |
full_conversation |
list[Message] |
用于链式处理的完整会话上下文(先前的输入 + 代理输出)。 |
链式代理执行器时,下游执行器通过 AgentExecutorResponse 处理程序接收 from_response。 默认情况下,它使用 full_conversation 字段来保留完整的会话历史记录,从而阻止下游代理丢失以前的上下文。 可以使用 上下文模式更改此行为:
spam_detector = AgentExecutor(create_spam_detector_agent())
email_assistant = AgentExecutor(create_email_assistant_agent())
# The email_assistant receives the spam_detector's full conversation context
workflow = (
WorkflowBuilder(start_executor=spam_detector)
.add_edge(spam_detector, email_assistant)
.build()
)
流式处理行为
AgentExecutor 自动适应工作流执行模式。
-
stream=True— 调用agent.run(stream=True)并生成每个AgentResponseUpdate事件作为工作流输出事件。 流式处理完成后,更新将聚合成完整的AgentResponse以供下游发送。 -
stream=False(默认值) - 调用agent.run(stream=False)并生成单个AgentResponse作为工作流输出事件。
# Streaming mode — receive incremental updates
events = workflow.run("Write a story about a cat.", stream=True)
async for event in events:
if event.type == "output" and isinstance(event.data, AgentResponseUpdate):
print(event.data.text, end="", flush=True)
# Non-streaming mode — receive complete response
result = await workflow.run("Write a story about a cat.")
# Retrieve terminal AgentResponse objects from the result
outputs = result.get_outputs()
for output in outputs:
if isinstance(output, AgentResponse):
print(output.text)
# Retrieve intermediate outputs (progress / observational emissions)
intermediate_outputs = result.get_intermediate_outputs()
for item in intermediate_outputs:
print(f"Intermediate: {item}")
上下文模式
当代理被链接在一起时,context_mode 上的 AgentExecutor 参数用于控制代理通过 AgentExecutorResponse 处理程序从上游代理接收 from_response 时所消耗的会话上下文。
可用模式
| 模式 | 行为 |
|---|---|
"full"(默认值) |
该代理使用上游代理的完整会话 - 提供给上游代理的输入消息及其响应消息。 |
"last_agent" |
代理仅使用上游代理的响应消息,不包括提供给上游代理的输入。 |
"custom" |
用户提供的 context_filter 函数确定代理使用的消息。 需要 context_filter 参数。 |
使用 last_agent 模式
使用"last_agent",当每个代理只专注于转换上一个代理的输出,而不受早期会话轮次的影响。 这对于转换管道、渐进式优化和类似的顺序转换非常有用:
from agent_framework import AgentExecutor, WorkflowBuilder
# Each agent consumes only the previous agent's response messages
french_executor = AgentExecutor(french_agent, context_mode="last_agent")
spanish_executor = AgentExecutor(spanish_agent, context_mode="last_agent")
workflow = (
WorkflowBuilder(start_executor=writer_agent)
.add_edge(writer_agent, french_executor)
.add_edge(french_executor, spanish_executor)
.build()
)
使用 context_mode="last_agent"时,法国翻译只使用编写器的响应消息(不包括输入给作者的原始用户提示),而西班牙语翻译只使用法语翻译的响应消息。
使用 custom 模式
若要对代理使用的上下文进行精细控制,请结合 context_mode="custom" 和 context_filter 函数一起使用。 筛选器接收完整对话作为list[Message],并返回经过筛选的子集:
from agent_framework import AgentExecutor, Message
def keep_user_and_last_agent(messages: list[Message]) -> list[Message]:
"""Keep only user messages and the last agent's response."""
user_msgs = [m for m in messages if m.role == "user"]
agent_msgs = [m for m in messages if m.role == "assistant"]
return user_msgs + agent_msgs[-1:] if agent_msgs else user_msgs
executor = AgentExecutor(
my_agent,
context_mode="custom",
context_filter=keep_user_and_last_agent,
)
SequentialBuilder 中的上下文模式
编排
from agent_framework.orchestrations import SequentialBuilder
workflow = SequentialBuilder(
participants=[writer, translator, reviewer],
chain_only_agent_responses=True,
).build()
有关完整示例,请参阅 Agent Framework 存储库中的 sequential_chain_only_agent_responses.py 。
共享会话
默认情况下,每个 AgentExecutor 会创建自己的会话。 若要在多个代理之间共享会话(例如,若要维护公共会话线程),请显式创建会话并将其传递给每个执行程序:
from agent_framework import AgentExecutor
# Create a shared session from one agent
shared_session = writer_agent.create_session()
# Both executors share the same session
writer_executor = AgentExecutor(writer_agent, session=shared_session)
reviewer_executor = AgentExecutor(reviewer_agent, session=shared_session)
注释
并非所有代理都支持共享会话。 通常,只有同一提供程序类型的代理可以共享会话。
检查点
在长时间运行的工作流中,AgentExecutor 支持检查点以保存和还原状态。 当进行检查点时,执行器会进行序列化:
- 内部消息缓存。
- 完整的对话历史记录。
- 代理会话状态。
- 任何挂起的用户输入请求和响应。
在还原时,执行程序会反序列化此状态,从而允许工作流从其离开位置恢复。
警告
使用服务器端会话(例如 FoundryAgent)的代理程序在进行检查点时存在限制。 服务器端会话状态未在检查点中捕获,可由后续运行修改。 如果需要使用服务器端会话进行可靠的检查点机制,请考虑实现自定义执行器。
工作原理
Go 使用 workflow/agentworkflow 将代理托管为工作流执行器。 托管执行程序使用以下 轮次令牌 模式:
- 消息缓冲 - 当消息从其他执行程序到达时,托管代理会收集它们。 如果启用消息转发(默认值),则传入消息也会转发到下游执行程序。
-
轮次令牌触发 — 托管代理仅在收到
workflow.TurnToken后才会处理其缓存的消息。 -
代理调用 — 执行器通过
Run调用底层代理,并从agentworkflow.Config或TurnToken中选择流式行为。 -
输出生成 - 如果启用了更新事件,则每个
*agent.ResponseUpdate事件作为工作流输出生成。 如果启用了响应事件,则会将聚合后的*agent.Response作为工作流输出。 - 下游消息传送 - 代理的响应消息将发送到连接的下游执行程序。
-
轮次令牌传递——当前轮次完成后,执行程序会向下游发送一个新的
workflow.TurnToken,以下一个托管代理开始处理。
自定义配置
通过使用 agentworkflow.New 和 agentworkflow.Config 值创建绑定,自定义托管代理执行器的行为:
hostedAgent := agentworkflow.New(myAgent, agentworkflow.Config{
EmitUpdateEvents: true,
DisableForwardIncomingMessages: true,
})
wf, err := workflow.NewBuilder(hostedAgent).
WithOutputFrom(hostedAgent).
Build()
if err != nil {
return err
}
小窍门
有关完整的可运行示例,请参阅 工作流示例中的代理 。
输入类型
托管代理执行器接受 string、*message.Message、[]*message.Message 和 iter.Seq[*message.Message] 输入。 字符串输入会被转换为具有 User 角色的 message.Message 实例。 消息输入会被缓冲,直到执行器收到 workflow.TurnToken,这会触发托管代理针对累计的这批消息运行。
run, err := inproc.Default.RunStreaming(ctx, wf, nil)
if err != nil {
return err
}
defer run.Close(ctx)
if err := run.SendMessage(ctx, "Summarize this deployment plan."); err != nil {
return err
}
if err := run.SendMessage(ctx, message.NewText("Include risk notes.")); err != nil {
return err
}
if err := run.SendMessage(ctx, []*message.Message{message.NewText("Keep it concise.")}); err != nil {
return err
}
emitEvents := true
if err := run.SendMessage(ctx, workflow.TurnToken{EmitEvents: &emitEvents}); err != nil {
return err
}
输出和链式调用
托管代理完成轮次后,它会将代理的响应消息和新轮次令牌发送到连接的下游执行程序。 这使得将代理串联起来变得非常简单:
french := agentworkflow.New(frenchAgent, agentworkflow.Config{})
spanish := agentworkflow.New(spanishAgent, agentworkflow.Config{})
english := agentworkflow.New(englishAgent, agentworkflow.Config{})
wf, err := workflow.NewBuilder(french).
AddEdge(french, spanish).
AddEdge(spanish, english).
Build()
if err != nil {
return err
}
流式处理行为
在 agentworkflow.Config 上设置 EmitUpdateEvents,或发送设置了 EmitEvents 的 workflow.TurnToken,以便通过工作流输出事件发出代理响应更新。
hostedAgent := agentworkflow.New(myAgent, agentworkflow.Config{
EmitUpdateEvents: true,
})
wf, err := workflow.NewBuilder(hostedAgent).
WithOutputFrom(hostedAgent).
Build()
if err != nil {
return err
}
run, err := inproc.Default.RunStreaming(ctx, wf, message.NewText("Write a status update."))
if err != nil {
return err
}
defer run.Close(ctx)
emitEvents := true
if err := run.SendMessage(ctx, workflow.TurnToken{EmitEvents: &emitEvents}); err != nil {
return err
}
for evt, err := range run.WatchStream(ctx) {
if err != nil {
return err
}
if output, ok := evt.(workflow.OutputEvent); ok {
if update, ok := output.Output.(*agent.ResponseUpdate); ok {
fmt.Print(update.String())
}
}
}
配置选项
agentworkflow.Config 控制托管代理执行程序的行为:
| 选项 | 默认 | 说明 |
|---|---|---|
EmitUpdateEvents |
false |
在执行期间输出流式 *agent.ResponseUpdate 值。
workflow.TurnToken.EmitEvents 在已设置时优先。 |
EmitResponseEvents |
false |
将聚合后的 *agent.Response 作为工作流输出事件发出。 |
InterceptUserInputRequests |
false |
截获 ToolApprovalRequestContent 并将其作为工作流消息进行路由处理。 |
InterceptUnterminatedFunctionCalls |
false |
截获未解析 FunctionCallContent 的值,并将其路由为工作流消息。 |
DisableReassignOtherAgentsAsUsers |
false |
保留来自其他代理的传入助理角色,而不是将其重新分配给用户角色。 |
DisableForwardIncomingMessages |
false |
在托管代理生成的消息之前,停止将传入消息转发到下游执行程序。 |
hostedAgent := agentworkflow.New(myAgent, agentworkflow.Config{
EmitUpdateEvents: true,
EmitResponseEvents: true,
InterceptUserInputRequests: true,
InterceptUnterminatedFunctionCalls: true,
DisableReassignOtherAgentsAsUsers: false,
DisableForwardIncomingMessages: false,
})
检查点
托管代理参与工作流检查点机制。
agentworkflow.New 在执行程序上注册检查点和还原挂钩。 执行检查点时,主机存储:
- 托管代理的
agent.SessionJSON 状态。 - 当前轮次的事件发送设置。
- 待处理的工具审批和函数调用请求状态
在还原时,主机会在工作流继续之前重新创建代理会话并还原挂起的请求处理程序。 通过工作流执行环境启用检查点功能,例如使用 inproc.Default.WithCheckpointing(...);不需要 agentworkflow.Config 选项。
checkpointManager := checkpoint.NewInMemoryManager()
environment := inproc.Default.WithCheckpointing(checkpointManager)
var checkpoints []workflow.CheckpointInfo
run, err := environment.RunStreaming(ctx, wf, message.NewText("Start the review."))
if err != nil {
return err
}
defer run.Close(ctx)
emitEvents := true
if err := run.SendMessage(ctx, workflow.TurnToken{EmitEvents: &emitEvents}); err != nil {
return err
}
for evt, err := range run.WatchUntilHalt(ctx) {
if err != nil {
return err
}
if completed, ok := evt.(workflow.SuperStepCompletedEvent); ok && completed.CompletionInfo != nil {
if completed.CompletionInfo.CheckpointInfo != nil {
checkpoints = append(checkpoints, *completed.CompletionInfo.CheckpointInfo)
}
}
}
if len(checkpoints) == 0 {
return fmt.Errorf("no checkpoints were created")
}
resumedRun, err := environment.ResumeStreaming(ctx, wf, checkpoints[len(checkpoints)-1])
if err != nil {
return err
}
defer resumedRun.Close(ctx)
注释
由提供程序支持的会话仍可能存在提供程序特有的持久性限制。 检查点会捕获 Go 主机可用的 agent.Session 状态,而不是提供程序未序列化到会话中的外部服务状态。