跳到主要内容

数据截至 (上游 commit 5e1f1fb87d9a)

多 agent 编排:actor 运行时上的五种协作模式

30 秒导读: 前面几章讲的是"一个 agent 怎么跑"(04)。这一章讲"多个 agent 怎么一起跑"。 SK 的答案不是写五套调度循环,而是先做一个通用 actor 运行时(只负责把消息送到该去的地方), 再在它上面用五种消息协议拼出五种协作形态:并发、串行、群聊、交接、Magentic。

成熟度提醒: 本章涉及的每一个类都带 @experimental 装饰器 (python/semantic_kernel/utils/feature_stage_decorator.pyexperimental)。 从 OrchestrationBaseInProcessRuntime 到五种编排,全部是"随时可能改签名"的实验特性。


1. 这是什么(零基础也能懂)

一句话定义: 多 agent 编排层 = 一个进程内的消息中间件(actor 运行时)+ 五个"协作剧本"模板。

它解决什么问题。 假设你有三个 agent:物理专家、化学专家、写作专家。你想让它们:

  • 同时回答同一个问题,最后把三份答案收齐 —— 这是并发;
  • 一个接一个加工同一份草稿 —— 这是串行;
  • 围着一个话题轮流发言,由某个"主持人"决定谁说话、什么时候散会 —— 这是群聊;
  • 客服 agent 发现自己搞不定,把对话转交给退款 agent —— 这是交接;
  • 一个"项目经理" agent 先写计划,再逐轮点人干活、检查进度、卡住就重新规划 —— 这是 Magentic

如果每种都手写一遍循环,你会重复实现五遍"谁给谁发消息、消息怎么排队、怎么取消、怎么收结果"。 SK 把这块公共的东西抽成了运行时,五个模板只写各自的消息协议

用起来什么样。 下面是仓库自带样例的骨架 (python/samples/getting_started_with_agents/multi_agent_orchestration/step1_concurrent.py):

# 1. 挑一个编排模板,把成员塞进去
orchestration = ConcurrentOrchestration(members=[physics_agent, chemistry_agent])

# 2. 起一个运行时(消息泵在后台任务里转)
runtime = InProcessRuntime()
runtime.start()

# 3. 投任务 —— 这一步立刻返回,不等结果
result = await orchestration.invoke(task="What is temperature?", runtime=runtime)

# 4. 想要结果的时候再等
value = await result.get(timeout=20)

# 5. 队列排空后停机
await runtime.stop_when_idle()

一句话直觉: 把它想成公司内部的邮件系统 + 五种会议流程。 运行时是邮件系统(只管投递,不管你信里写了啥);编排模板是会议流程 (全员各写一份 / 击鼓传花 / 主持人点名 / 谁接手谁负责 / 项目经理带着计划推进)。


2. 顶层全景(它大概怎么转)

2.1 三层结构

这张图从下往上读:每一层只依赖它下面那层,越往上越"懂业务"。

第 3 层 五种编排模板 concurrent | sequential | group_chat | handoffs | magentic
管什么 定义消息类型、谁给谁发、什么时候算结束

│ 复用
第 2 层 粘合层 AgentActorBase
管什么 收到消息 → 攒进 message_cache → 调 agent.invoke_stream()
→ 把完整回答再发出去;顺带做异常上报和流式回调

│ 复用
第 1 层 actor 运行时 InProcessRuntime
管什么 AgentId 寻址、TopicId 广播、订阅匹配、单队列消息泵
不管什么 完全不知道什么是 LLM、什么是 prompt

2.2 部件一句话职责

部件干什么在哪个文件
OrchestrationBase编排的统一外壳:invoke 立即返回、准备 actor、做输入输出类型转换python/semantic_kernel/agents/orchestration/orchestration_base.py:85
OrchestrationResult结果句柄:get() 等一个 asyncio.Eventcancel() 中途叫停同上 :34
InProcessRuntime单 asyncio 队列的消息泵,负责投递、懒实例化 actor、存取状态python/semantic_kernel/agents/runtime/in_process/in_process_runtime.py:167
RoutedAgentactor 基类:按消息类型把消息分发到对应的 handler 方法python/semantic_kernel/agents/runtime/core/routed_agent.py:447
TypeSubscription订阅规则:话题类型对上就转成 AgentId,话题的 source 当实例 keypython/semantic_kernel/agents/runtime/in_process/type_subscription.py:14
AgentActorBase把一个 SK Agent 装进 actor 壳里的粘合层python/semantic_kernel/agents/orchestration/agent_actor_base.py:97
五个编排模板各自的消息协议 + 终止条件orchestration/{concurrent,sequential,group_chat,handoffs,magentic}.py

2.3 主线走一遍(不进代码)

你调 invoke(task, runtime)

├─ ① 给这次运行造一个唯一 id(uuid4().hex),后面所有 actor 类型名和话题名都带上它
├─ ② _prepare():把每个 Agent 包成 actor 注册进运行时,再挂上订阅
├─ ③ 把 task 转成 ChatMessageContent(必要时走 input_transform)
├─ ④ 起一个 asyncio 后台任务跑 _start():发出第一条消息
└─ ⑤ 立刻 return OrchestrationResult(此刻还是空的)


运行时的消息泵开始转 → actor 们互相收发 → 某个 actor 拿到最终结果

└─ 调 result_callback → 走 output_transform → 填值 → event.set()

你在别处 await result.get() ◀────────────────────────────────────────────┘ 被唤醒

3. 编排契约:invoke 为什么"立刻返回"

它要解决的小问题。 一次多 agent 协作可能跑几分钟。如果 invoke 一路阻塞到结束, 调用方就没法在中途取消、也没法一边等一边干别的。

思路。 把"启动"和"取结果"拆成两个动作:invoke 只负责点火并返回一个句柄, 结果通过 asyncio.Event 交付。

3.1 结果句柄:一个 Event 就够了

OrchestrationResult 里只有五个字段,核心是 event (orchestration_base.py:34-41)。等结果的逻辑在 get(:43-67):

if timeout is not None:
await asyncio.wait_for(self.event.wait(), timeout=timeout)
else:
await self.event.wait()

超时抛 TimeoutError不中止运行(docstring 明说了,:47-49)。 get 醒来后如果 value 还是 None,会按"被取消 / 有异常 / 无结果"三种情况分别报错(:61-66)。

cancel()(:69-81)干两件事:置 CancellationToken、把 event 设起来让 get 立刻醒。 注意它不会清空运行时队列——docstring 写得很直白:已经收到消息的 actor 会继续处理完。 真正的"刹车"在下游:ActorBase.on_message_impl 每次进来先看令牌,取消了就直接 return None (agent_actor_base.py:45-53)。

3.2 每次运行一个 uuid:同一个运行时能跑多场

这是整套设计里最省事的一招(orchestration_base.py:214-215):

# This unique topic type is used to isolate the orchestration run from others.
internal_topic_type = uuid.uuid4().hex

这个 hex 串同时被用作话题类型actor 类型名后缀。各编排里都有一模一样的一行:

return f"{agent.name}_{internal_topic_type}" # concurrent.py:239 / sequential.py:197 /
# group_chat.py:519 / handoffs.py:512 / magentic.py:900

带来两个结果:

  • 同一个 InProcessRuntime 上可以并行跑好几场编排,话题互不串台;
  • 同一个编排对象 invoke 两次也不会撞上 register_factory 的"类型已存在"报错 (in_process_runtime.py:740-741)。

3.3 输入输出转换:让编排接受任意类型

运行时内部只认 ChatMessageContent(或它的 list),别名叫 DefaultTypeAlias (orchestration_base.py:27)。但用户可能想传一个 pydantic 模型进去、拿一个 pydantic 模型出来。

invoke 的处理顺序是"能直接用就直接用,不行才转"(:224-234):

传进来的 task怎么处理
str包成 ChatMessageContent(role=USER)
ChatMessageContent 或它的 list原样透传
其它类型交给 input_transform(可以是 async)

默认转换 _default_input_transform(:295-319)对自定义类型做的是 json.dumps(input_message.__dict__) ——简单粗暴,把对象序列化成一段 JSON 文本喂给 agent。 反方向 _default_output_transform(:321-344)则是 self.t_out(**json.loads(output_message.content)), 指望模型把 JSON 原样吐回来。

想要更靠谱的结构化输出,仓库另给了 structured_outputs_transform (python/semantic_kernel/agents/orchestration/tools.py:18):它另起一个 chat completion 调用, 用 response_format 把自由文本压成目标 schema。

类型参数本身是靠 _set_types()__orig_class__ / __orig_bases__ 里挖出来的(:128-179), 所以 ConcurrentOrchestration[MyIn, MyOut](...) 这种写法才能工作。


4. actor 运行时(第 1 层,核心机制)

这一节讲的东西完全不涉及 LLM。你可以把 python/semantic_kernel/agents/runtime/ 整个目录 当成一个独立的迷你 actor 框架来读。

4.1 三个概念:身份、话题、订阅

概念是什么字符串形态定义处
AgentId一个 actor 实例的地址 = 类型 + 实例 keytype/keyruntime/core/agent_id.py:19(协议)、:49 CoreAgentId
AgentTypeactor 的类别(一个类型可以有多个实例)typeruntime/core/agent_type.py:12
TopicId广播的作用域 = 事件类型 + 来源type/sourceruntime/core/topic.py:19
Subscription规则:哪些 TopicId 该落到哪个 AgentIdruntime/core/subscription.py:13

订阅只有两个实现,逻辑都短得一眼看完:

# TypeSubscription:话题类型完全相等
def is_match(self, topic_id): return topic_id.type == self._topic_type # type_subscription.py:51-53
def map_to_agent(self, topic_id): return CoreAgentId(self._agent_type, topic_id.source) # :55-60
# TypePrefixSubscription:话题类型前缀匹配
def is_match(self, topic_id): return topic_id.type.startswith(self._topic_type_prefix) # type_prefix_subscription.py:49-51

这里有个关键设计:话题的 source 被拿去当 actor 的实例 key。 所以同一个订阅规则可以按 source 分裂出多个独立实例——"每个会话一份 actor"。 五种编排都用同一个 source 值,所以实际上每场运行每种 actor 只有一个实例。

SubscriptionManager 做的是懒缓存:第一次见到某个 topic 就遍历所有订阅算一遍收件人, 之后直接查表;加/删订阅时把见过的话题全部重算 (runtime/in_process/runtime_impl_helpers.py:76-94)。

还有一条容易漏掉的:BaseAgent.register 除了注册工厂,还会自动给每个 actor 类型加一条 前缀订阅 TypePrefixSubscription(agent_type + ":"),专门用来接"点对点的广播" (runtime/core/base_agent.py:207-215,注释明说前缀必须带 : 防撞名)。

4.2 三种信封 + 一条队列

运行时里流动的不是裸消息,而是三种信封(dataclass):

信封什么时候产生关键字段
SendMessageEnvelopesend_message(点对点,要回复)recipientfuture
PublishMessageEnvelopepublish_message(广播,不要回复)topic_id
ResponseMessageEnvelope处理完 send 后回程futuresender/recipient

依据:in_process_runtime.py:68(publish)、:81(send)、:95(response)。

三种信封挤在同一条 asyncio.Queue 里(:183)。消息泵长这样:

send_message ──┐ ┌─▶ _process_send ──┐
├─▶ [ 一条 Queue ] ──▶ _process_next ──┤ │
publish_message┘ FIFO (match 信封类型) ├─▶ _process_publish │
└─▶ _process_response

回程信封重新入队 ◀───────────────────┘

怎么读这张图: 所有消息先排队,_process_next 取一个、按类型分派, 每个分派出去的处理都是 asyncio.create_task —— 所以"投递有序、处理并发"。

几处值得看的实现:

  • send 是同步语义。 send_message 建一个 future,塞信封,然后 await future(:241-262)。 被 cancellation_token.link_future(future) 挂上取消钩子(:260)。
  • 回程要绕一圈队列。 _process_send 拿到返回值后不直接 set_result, 而是再包一个 ResponseMessageEnvelope 入队(:417-425),由 _process_response 兑现(:516-517)。 好处是拦截器(InterventionHandler)对回程消息也有一次拦截机会(:623-644)。
  • publish 不回发给自己。 _process_publish 里显式跳过 sender(:434-436)。 这条不起眼的规则直接影响上层:Magentic 的 manager 必须手动把自己发的消息补进自己的历史 (magentic.py:572-582 的注释就是在解释这件事)。
  • 单次分派后让出控制权。 _process_next 末尾一句 await asyncio.sleep(0)(:650), 保证消息泵不会独占事件循环。
  • 异常策略可选。 构造参数 ignore_unhandled_exceptions 默认 True; 设为 False 时 publish 里的异常会被存进 _background_exception,在下一次 _process_next 开头抛出(:489-491:530-534)。

停机有三档,日常只用第二个:

方法行为位置
stop()立刻停,队列里剩下的消息全丢:668-679
stop_when_idle()queue.join() 排空再停,最常用:682-694,实现在 RunContext.stop_when_idle(:132-137)
stop_when(cond)轮询条件,源码自己标了 .. caution:: 不推荐:696-712

Python < 3.13 的 asyncio.Queue 没有 shutdown,所以仓库自带了一份回填实现—— runtime/in_process/queue.py 整份文件都是 CPython 队列源码的拷贝。 按版本二选一导入的那四行不在 queue.py 里,而在 in_process_runtime.py:16-19: 3.13 及以上取标准库的 Queue/QueueShutDown,否则取这份回填版。

4.3 actor 懒实例化与状态

register_factory 只登记一个工厂函数,不创建对象(:729-757)。 真正创建发生在第一次有消息要投给某个 AgentId 时,由 _get_agent 触发(:795-805):

消息要发给 AgentId("Writer_ab12", "default")

├─ 已经在 _instantiated_agents 里? ──是──▶ 直接用
└─ 否 ──▶ 查 _agent_factories[type] ──▶ 在 AgentInstantiationContext 里调工厂 ──▶ 缓存

那个上下文变量是必需的:BaseAgent.__init__ 会从中取出 runtime 和自己的 id, 拿不到就直接抛"不能脱离运行时直接构造"(base_agent.py:96-103)。

save_state / load_state(:309-342)遍历已实例化的 actor,按 str(agent_id) 存成字典。 两处 docstring 都注明了:订阅状态目前不存(:315-316:334-336)。

4.4 RoutedAgent:按类型路由到方法

它要解决的小问题。 一个 actor 要处理好几种消息(开始、请求、回复、重置)。 不想写一长串 isinstance 分支。

思路。 用装饰器把每个方法的消息参数类型注解读出来,建一张 类型 → handler 列表 的表。

RoutedAgent.__init__ 扫描类里所有带 is_message_handler 标记的方法建表 (routed_agent.py:463-473_discover_handlers:502-511)。 分发只有五行(:475-489):

handlers = self._handlers.get(type(message))
if handlers is not None:
for h in handlers:
if h.router(message, ctx):
return await h(self, message, ctx)
return await self.on_unhandled_message(message, ctx)

注意是 type(message) 精确匹配,不走继承(源码 :46 的注释原话就是"works on concrete types and not inheritance")。

三个装饰器的差别只在那个 router lambda 上,一行定生死:

装饰器router 逻辑效果位置
@message_handlermatch or (lambda: True)广播和 RPC 都收:169
@event(not ctx.is_rpc) and match只收广播,且必须返回 None:301
@rpcctx.is_rpc and match只收点对点:429

ctx.is_rpc 不是 routed_agent.py 自己填的,而是运行时按信封种类填的: send 那条路填 True(in_process_runtime.py:371),publish 那条路填 False(同文件 :458)。

有意思的是:五种编排一个 @event、一个 @rpc 都没用,全部用最宽的 @message_handler。 所以同一个 handler 既能接 send_message 也能接 publish_message —— sequential 的 CollectionActor 正是靠这一点接收上一环 send 过来的消息(sequential.py:98-101)。

装饰器还顺手做了运行时类型校验:收到不在 target_types 里的消息, strict=True(默认)就抛 CantHandleException(:151-154);返回值类型不对也抛(:158-161)。


5. 粘合层:把一个 Agent 装进 actor 壳

AgentActorBase(agent_actor_base.py:97)是第 3 层所有 agent actor 的父类。 它的职责可以用一张图说完:

收到消息

├─▶ 攒进 _message_cache(一个 ChatHistory,只是暂存区)

└─▶ 轮到自己发言时调 _invoke_agent()

├─ _create_messages():把 cache 里的消息取出来 + 清空 cache
├─ agent.invoke_stream(messages, thread=self._agent_thread, ...)
├─ 边收 chunk 边回调 streaming_agent_response_callback
├─ 把 chunk 加起来拼成完整回答(sum(...))
└─ 回调 agent_response_callback,返回完整回答

几个真实细节:

  • cache 用完即清。 _create_messages 先复制再 clear()(:210-213)。 为什么可以清?因为对话连续性由 AgentThread 承担——第一次调用后 self._agent_thread = response_item.thread 就把线程接住了(:185-186)。 线程是什么见 04 章
  • 只走流式接口。 无论用户要不要流式,_invoke_agent 一律用 invoke_stream, 再用 sum(streaming_message_buffer[1:], streaming_message_buffer[0]) 拼回完整消息(:196)。 is_final 标志靠"缓冲一个 chunk 再回调"实现——回调的永远是上一个 chunk(:183-190)。
  • 一个都没收到就报错。 raise RuntimeError(f'Agent "{...}" did not return any response.')(:193)。 handoff 编排必须绕开这条,后面 §6.4 会讲。
  • 异常有统一出口。 ActorBase.exception_handler 是个装饰器(:56-93), 捕获后先调 exception_callback(最终填进 OrchestrationResult.exception)再重新抛出。 同步/异步两种包装都写了。

6. 五种编排:各自的消息协议与终止条件

先看总表,再逐个展开。

编排谁在指挥消息怎么流什么时候停结果是什么
concurrent没人(纯扇出)一次 publish 扇出,各自 send 回收集器收集器数满 N 份list[ChatMessageContent]
sequential没人(链在注册时写死)send 一环扣一环走到链尾的收集器最后一环的回答
group_chatGroupChatManager主持人 publish 点名 → 被点名者 publish 回答should_terminate 为真filter_results 挑出的一条
handoffsagent 自己(靠函数调用)publish 点名 → 交接或作答有人调用 complete_task任务总结文本
magenticStandardMagenticManager主持人 publish 指令+点名 → 被点名者回答进度账本判定满足,或撞上限prepare_final_answer 的产物

6.1 concurrent:扇出扇入

要解决的问题: 同一个任务问 N 个专家,收齐 N 份答案。

_start: publish(ConcurrentRequestMessage, Topic(uuid, "ConcurrentOrchestration"))
│ 订阅匹配,一次扇出到全部成员
┌───────┼───────┐
▼ ▼ ▼
ActorA ActorB ActorC 各自 _invoke_agent()
│ │ │
└───────┼───────┘ send(ConcurrentResponseMessage) 点对点

CollectionActor 数够 N 份 → result_callback(list)
  • 扇出:_start 只发一条 publish(concurrent.py:143-147),订阅在 _add_subscriptions 里批量加好(:218-231)。
  • 扇入:每个 actor 处理完用 send_message 直接寄给收集器(:92-97),不是广播——收集器不订阅话题。
  • 终止:CollectionActor 拿锁 append,然后比较 len(self._results) == self._expected_answer_count(:119-127)。
  • 顺序不保证,样例注释也写了("the order of the results is not guaranteed")。

6.2 sequential:链式传递

要解决的问题: 一份内容逐道工序加工。

_start: send ──▶ Actor1 ──send──▶ Actor2 ──send──▶ Actor3 ──send──▶ CollectionActor
│ │ │ │
回答变成下一环的 … … result_callback
SequentialRequestMessage

这里最巧的是注册顺序反着来(sequential.py:156-171):

next_actor_type = self._get_collection_actor_type(internal_topic_type)
for agent in reversed(self._members):
await SequentialAgentActor.register(runtime, ..., lambda ... next_agent_type=next_actor_type: ...)
next_actor_type = self._get_agent_actor_type(agent, internal_topic_type)

为什么必须反着注册: 每个 actor 在构造时就要知道"下一环是谁"。 从尾往头注册,轮到某个 agent 时它的后继类型名已经确定了。docstring 把这个理由写在了 :142-155

注意链条是注册期写死的静态链,运行期不能改道。整条链上没有订阅,全靠 send_message

6.3 group_chat:主持人点名

要解决的问题: 让一组 agent 围绕一个话题轮流发言,并且能插入人类。

先和另一套"群聊"划清界限: 04 章 §9 讲的 AgentGroupChat / AgentChannel / BroadcastQueue与本章并存的旧版群聊路径——它让异构 agent 共享一段 权威历史,跑在一把进程内锁上。本节的 group_chat 编排跑在 actor 运行时上,agent 之间只传消息、 不共享历史,也不需要 channel。两者名字像,机制无关,不要混着读。

消息只有三种(group_chat.py:40/47/54):GroupChatStartMessageGroupChatRequestMessage(agent_name)GroupChatResponseMessage(body)

_start ──send──▶ 每个成员 actor(先灌入初始任务,再灌主持人)
⚠ 顺序重要,见下文
主持人 ──publish(RequestMessage{agent_name:"Writer"})──▶ 全体
│ 每个 actor 自查:不是叫我 → return

Writer 作答 ──publish(ResponseMessage)──▶ 全体(含主持人)


主持人重新走一遍决策 → 继续点名 或 收工

点名靠自筛。 请求是广播的,每个 actor 开头一句 if message.agent_name != self._agent.name: return(:123-125)。

主持人的四个决策钩子(GroupChatManager,:147),每收到一条回答就按顺序跑一遍 (_determine_state_and_take_action,:307-350):

钩子问什么抽象?位置
should_request_user_input现在要不要问人类抽象:156
should_terminate该散会了吗有默认实现:164-179
select_next_agent下一个谁说抽象:181-193
filter_results从历史里挑哪条当最终结果抽象:195-206

should_terminate 的默认实现就是轮数守卫:先 self.current_round += 1, 再返回 current_round > max_rounds(:170-178)。子类若要加自己的判断,通常得先 await super()

内置的 RoundRobinGroupChatManager(:210-243)三个抽象方法都实现得极简: 不要人类输入、按 current_index 取模轮转、结果取历史最后一条。

两个真实细节:

  • 每次调钩子传的都是 self._chat_history.model_copy(deep=True)(:311:326:330:337)—— 深拷贝,防止 manager 实现顺手改坏共享历史。
  • 终止时把两个理由塞进结果的 metadata:termination_reasonfilter_result_reason(:331-332)。

_start 为什么要先给成员发再给主持人发? 源码 docstring 直接解释了(:420-426): 如果主持人太快、成员太慢,主持人可能在别人还没拿到任务上下文时就点了名。 所以 _startasyncio.gather 给所有成员 send(同步语义,等确认),最后才给主持人发(:429-445)。

6.4 handoffs:把"转交"编译成一个函数

要解决的问题: 让 agent 自己决定把对话交给谁,而不是由外部主持人调度。

思路: 既然 agent 本来就会调函数(见 03 章), 那就把"可以交接给谁"编译成一组 kernel function 塞进它的 kernel。模型选哪个函数 = 选交给谁。

和 04 章那把"工具柄"的关系: 04 章 §4 说每个 Agent 出生自带一个 _as_kernel_function 外壳。handoff 没有复用它,而是自己另造了一批 transfer_to_* 函数。 同一个"agent 变函数"的思路,仓库里有两处独立实现。

交接图用 OrchestrationHandoffs(handoffs.py:57)声明,本质是 dict[str, dict[str, str]]: 源 agent → {目标 agent: 这条路的描述}。add / add_many 支持链式调用(:76-117)。 构造时会校验:目标必须是成员、不能自交接(_validate_handoffs,:514-527)。

编译过程(_add_handoff_functions,:190-222):

handoff_connections = {"RefundAgent": "转给退款专员处理退货"}

├─▶ 造一个 KernelFunctionMetadata,name = "transfer_to_RefundAgent"
│ description = 那条路的描述 ← 模型就是靠这句话选路的
│ 用 partial(self._handoff_to_agent, "RefundAgent") 当方法体

├─▶ 再加一个真 @kernel_function:complete_task(task_summary)

├─▶ 全部塞进插件 "Handoff",加到 agent.kernel.clone() 上

└─▶ 给这个 clone 挂一个自动函数调用过滤器(AUTO_FUNCTION_INVOCATION)

两个关键点:

  1. kernel 是克隆的(:175 self._kernel = agent.kernel.clone()), 调用时用 kernel=self._kernel 覆盖(:275)。原 agent 的 kernel 不被污染。
  2. 过滤器负责刹车(:229-233):
async def _handoff_function_filter(self, context: AutoFunctionInvocationContext, next):
await next(context)
if context.function.plugin_name == HANDOFF_PLUGIN_NAME:
context.terminate = True

先让函数真正执行(_handoff_to_agentself._handoff_agent_name 记下来), 再置 context.terminate = True 掐断自动函数调用循环——所以模型不会再生成一段文本回答。

这就带来一个副作用:agent 可能一句话都没回。于是 handoff 专门写了 _invoke_agent_with_potentially_no_response(:321-362),和 _invoke_agent 只差最后一步: 没收到任何 chunk 时 return None 而不是抛错(:355-356),docstring 把理由写清楚了(:325-334)。

主循环(_handle_request_message,:268-312)的判定顺序:

被点名 → 调 agent(带 Handoff 插件的 kernel)

├─ 设了 handoff 目标? ──是──▶ publish(RequestMessage{目标}) → 交棒,break

├─ 没设,也没回答? ──▶ RuntimeError

└─ 有回答 ──▶ publish(ResponseMessage) 让全场看到

├─ 配了 human_response_function? ──是──▶ 取人类输入 → 再问一轮(循环继续)
└─ 否 ──▶ _complete_task("No handoff agent name provided ...") → 结束

_complete_task(:235-249)是唯一的正常出口:调 result_callback 交付一条 "Task is completed with summary: ...",并置 _task_completed = True

6.5 magentic:任务账本 + 进度账本的内外双循环

要解决的问题: 开放式的复杂任务——需要先规划、逐步推进、检测原地打转、必要时推倒重来。

这是五种里最重的一种。核心是两本"账本":

账本内容什么时候写类型
task ledger(任务账本)facts 事实 + plan 计划开局一次;卡住时重写_TaskLedger(magentic.py:197)
progress ledger(进度账本)5 个判断项,见下每一轮都重新问模型ProgressLedger(:92-99)

进度账本的五项(每项都是 {reason, answer}):

字段问模型什么
is_request_satisfied任务完成了吗
is_in_loop是不是在原地打转
is_progress_being_made上一轮有进展吗
next_speaker下一个该谁
instruction_or_question给它的具体指令

它靠结构化输出保证能解析:create_progress_ledger 克隆一份执行设置, 把 response_format 设成 ProgressLedger 类本身,再 model_validate_json 回来(:450-463)。 StandardMagenticManager 构造时就校验服务支不支持结构化输出,不支持直接报错(:249-257)。

内外双循环。 注意这不是 while 语句,而是事件驱动的往返:

外层 _run_outer_loop (:568)
① 把渲染好的 task ledger publish 给全场(顺手补进自己的历史,因为 publish 不回发给自己)
② 调 _run_inner_loop

内层 _run_inner_loop (:596)
① _check_within_limits() → 超限就交付部分结果、return
② round_count += 1,问模型要一份 progress ledger
③ is_request_satisfied? ──是──▶ _prepare_final_answer() 收工
④ 没进展 或 在打转? ──▶ stall_count += 1;否则 stall_count 减 1(不低于 0)
⑤ stall_count > max_stall_count(默认 3)? ──是──▶ replan + 发 ResetMessage + 回到外层
⑥ 都不是 ──▶ publish 指令 + publish 点名 → 本次返回


被点名的 agent 作答 → publish(MagenticResponseMessage)

└─▶ manager 的 _handle_response_message (:549) 再次进入内层循环

三个细节:

  • 重置是真的清空。 MagenticResetMessage 到达 agent actor 后,清 _message_cache, 并且 await self._agent_thread.delete() 把线程删掉、置 None(:748-755)。 下一轮相当于全新对话。manager 自己的 MagenticContext.reset() 清历史、清 stall、reset_count += 1(:117-125)。
  • 两道上限守卫。 _check_within_limits(:681-714)看 max_round_countmax_reset_count。 撞上限时不是抛异常,而是尽力交付:从历史里倒着找最近一条 assistant 消息当部分结果; 一条都没有才编一句"Stopped because the maximum ... limit was reached"(:697-710)。
  • 提示词注入防护。 所有参与者描述在进模板前都过一遍 html.escape (:298-301:360-363:401-404:433-436),且渲染任务账本的模板显式 allow_dangerously_set_content=True(:398)——即"我知道这里放的是不可信内容,已自行转义"。 模板本身在 orchestration/prompts/_magentic_prompts.py

7. 巧妙之处(可借鉴的技术)

① 一个 uuid 同时当命名空间和话题名。 不需要"编排实例注册表", 拼在字符串里就实现了运行隔离(orchestration_base.py:215 + 五个 _get_agent_actor_type)。

② router lambda 把三个装饰器压成一个实现。 @event@rpc 的全部差别只是 not ctx.is_rpcctx.is_rpc 一个取反(routed_agent.py:301 vs :429)。

③ 用过滤器当"控制流劫持点"。 handoff 不需要改 agent 循环,只是往克隆的 kernel 上挂一个 自动函数调用过滤器,靠 context.terminate = True 掐断循环(handoffs.py:229-233)。 这是 03 章那套过滤器机制在多 agent 场景的复用。

④ 广播 + 自筛,而不是维护路由表。 群聊/交接/Magentic 的"点名"都是广播 agent_name, 每个 actor 自己判断是不是叫我(group_chat.py:124handoffs.py:271magentic.py:733)。 加减成员不用改任何路由逻辑。

⑤ 反向注册解决前向引用。 sequential 从链尾往链头注册,天然满足"每个节点构造时后继已知" (sequential.py:156-171)。

⑥ 撞上限时尽力返回部分结果,而不是抛异常让调用方一无所获(magentic.py:697-710)。

⑦ 决策钩子一律收深拷贝。 群聊每次调 manager 都传 model_copy(deep=True), 把"manager 实现写坏共享状态"这个 bug 类别从设计上排除(group_chat.py:311/326/330/337)。


8. 边界与局限

  • 全部 @experimental 没有稳定性承诺。
  • 只有进程内运行时。 runtime/ 下唯一实现是 InProcessRuntime; CoreRuntime 是个 Protocol(runtime/core/core_runtime.py:23),但仓库里没有跨进程/分布式实现。 消息序列化基建(serialization.py)已经在了,主要用于日志和 handles 声明。
  • 状态保存不完整。 save_state / load_state 的 docstring 自己写着不含订阅状态 (in_process_runtime.py:315-316:334-336)。而且只覆盖已实例化的 actor。
  • 取消是软取消。 OrchestrationResult.cancel() 不清空队列;已入队消息照样投递, 只是 actor 进门看到令牌就 no-op(agent_actor_base.py:50-51)。正在跑的那次 LLM 调用不会中断。
  • 消息路由不认继承。 RoutedAgenttype(message) 精确匹配(routed_agent.py:481), 给消息类型做子类不会命中父类的 handler。
  • 默认的输入/输出转换很朴素。 自定义类型走 json.dumps(obj.__dict__) 进、 t_out(**json.loads(content)) 出(orchestration_base.py:316:342), 完全依赖模型原样吐 JSON。要可靠就得用 structured_outputs_transform
  • concurrent 的收满判断在锁外。 append 在 async with self._lock 内, len(...) == expected 的比较在锁外(concurrent.py:121-124)。单事件循环下两句之间没有 await 点, 但这个不变量依赖于运行时是单线程的 (inferred)。
  • Magentic 依赖结构化输出。 模型不支持 response_format 就用不了标准 manager(magentic.py:249-257)。
  • handoff 无人类介入时只走一步。 没配 human_response_function 且 agent 也没交接, 就直接 _complete_task 收场(handoffs.py:308-312)。

9. 和本项目其它章节的关系

想知道什么去哪章
actor 里被调用的那个 Agent 本身怎么实现、AgentThread 是什么;以及与本章并存的旧版群聊 AgentGroupChat / AgentChannel / BroadcastQueue04-agents-and-threads.md(旧版群聊在该章 §9)
handoff 用的那个自动函数调用过滤器属于哪套机制03-function-calling-and-filters.md
Kernel / KernelFunction / 插件(handoff 把交接编译成的东西)01-kernel-and-functions.md
ChatMessageContent / 提示词模板渲染(Magentic 的账本模板)02-prompt-and-content.md
另一套编排:流程框架(Process Framework),以及数据检索层06-process-and-data.md
全局架构与阅读顺序index.md

和 Process 框架的区别(见 06): 本章的编排是运行期由消息驱动的松耦合协作; Process 框架是预先声明的有向图。想要确定性流水线用 06,想要 agent 自己商量着办用本章。


10. 代码地图(导航索引)

主题文件路径(相对克隆根)符号名
编排契约与结果句柄python/semantic_kernel/agents/orchestration/orchestration_base.pyOrchestrationBaseOrchestrationResultinvoke_prepare_start
输入输出类型转换同上_default_input_transform_default_output_transform_set_types
结构化输出转换器python/semantic_kernel/agents/orchestration/tools.pystructured_outputs_transform
Agent↔actor 粘合层python/semantic_kernel/agents/orchestration/agent_actor_base.pyActorBaseAgentActorBase_invoke_agentexception_handler
并发编排python/semantic_kernel/agents/orchestration/concurrent.pyConcurrentOrchestrationConcurrentAgentActorCollectionActor
串行编排python/semantic_kernel/agents/orchestration/sequential.pySequentialOrchestrationSequentialAgentActor_register_members
群聊编排python/semantic_kernel/agents/orchestration/group_chat.pyGroupChatManagerRoundRobinGroupChatManagerGroupChatManagerActor_determine_state_and_take_action
交接编排python/semantic_kernel/agents/orchestration/handoffs.pyOrchestrationHandoffsHandoffAgentActor_add_handoff_functions_handoff_function_filter_complete_task
Magentic 编排python/semantic_kernel/agents/orchestration/magentic.pyStandardMagenticManagerMagenticContextProgressLedger_run_outer_loop_run_inner_loop_check_within_limits
Magentic 提示词python/semantic_kernel/agents/orchestration/prompts/_magentic_prompts.pyORCHESTRATOR_TASK_LEDGER_*ORCHESTRATOR_PROGRESS_LEDGER_PROMPT
运行时接口python/semantic_kernel/agents/runtime/core/core_runtime.pyCoreRuntime
进程内运行时python/semantic_kernel/agents/runtime/in_process/in_process_runtime.pyInProcessRuntime_process_nextregister_factory_get_agentstop_when_idleRunContext
三种信封同上SendMessageEnvelopePublishMessageEnvelopeResponseMessageEnvelope
3.13 以下的队列回填python/semantic_kernel/agents/runtime/in_process/queue.pyQueueQueueShutDown(条件导入在 in_process_runtime.py:16-19)
类型路由python/semantic_kernel/agents/runtime/core/routed_agent.pyRoutedAgentmessage_handlereventrpcon_message_impl
actor 基类与注册python/semantic_kernel/agents/runtime/core/base_agent.pyBaseAgentregistersend_messagepublish_message
寻址与话题python/semantic_kernel/agents/runtime/core/{agent_id,agent_type,topic}.pyCoreAgentIdCoreAgentTypeTopicId
订阅匹配python/semantic_kernel/agents/runtime/in_process/type_subscription.pytype_prefix_subscription.pyTypeSubscriptionTypePrefixSubscriptionis_matchmap_to_agent
订阅缓存python/semantic_kernel/agents/runtime/in_process/runtime_impl_helpers.pySubscriptionManagerget_subscribed_recipients
取消与上下文python/semantic_kernel/agents/runtime/core/{cancellation_token,message_context}.pyCancellationTokenMessageContext
可运行样例python/samples/getting_started_with_agents/multi_agent_orchestration/step1_concurrent.pystep5_magentic.py