DCChains 动态编排链¶
DCChains 通过 LLM 生成 tf_plan 脚本来动态组合多个 Chain/Chains,适合处理包含分支、循环或重试逻辑的复杂任务。与 SeqChains 的固定顺序、ParallelChains 的并发执行相比,DCChains 把编排权交给可插拔的 orchestrate_chain 逻辑,使链路在运行时按需求重构。
核心特性¶
- 定义位置:
tfrobot.brain.chain.chain_structures.dc_chains.DCChains - 继承关系:继承
Chains,保持run/async_run接口 - 动态规划:
orchestrate_chain返回tf_plan脚本,由TFChainInterpreter解析并执行 - 结果聚合:所有子链的
ChainResult累积到单一返回值,保持上下文与 Token 统计
主要配置¶
| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
name |
ClassVar[str] |
链类型标识 | "DCChains" |
description |
ClassVar[str] |
描述信息 | 动态编排 Chains 的链结构 |
abilities_purpose |
Optional[str] |
提供给 LLM 的能力说明 | “当前的思维链是DCChains...” |
orchestrate_chain |
Callable[[ChainContext, DCChains, list[str], list[str]], str] |
负责编排的回调,返回包裹 tf_plan 的字符串 |
default_orchestrate_chain |
max_plan_times |
int |
规划失败后最多重试次数 | 3 |
orchestrate_chain 回调¶
- 入参包含:当前链上下文、DCChains 实例、历史失败的
tf_plan列表与错误信息。实现可以基于这些信息决定是否重新规划或降级策略。 - 返回值必须是形如
tf_plan\n...\n的字符串,运行前 DCChains 会去除外层标记并交给Lark+TFASTTransformer解析。 - 默认实现
default_orchestrate_chain使用 DeepSeek 模型生成脚本,自动把可用子链的输入/输出 schema 注入到模板中。
运行流程¶
- 上下文构造:
_run/_async_run首先把传入的current_input、conversation、knowledge、tools、intermediate_msgs封装到ChainContext。-
初始化空的
ChainResult,用于累计执行结果与 Token 统计。 -
规划与重试:
- 最多
max_plan_times次循环调用orchestrate_chain获取最新脚本。 -
每次失败都会把原脚本与失败原因写入
failed_plans/failed_infos,供下一轮规划参考。 -
脚本解析执行:
- 使用
Lark(PYTHON_GRAMMAR)解析脚本,TFASTTransformer转换为 AST。 -
TFChainInterpreter负责遍历 AST:根据脚本调用chains_map内注册的子链、执行工具、记录日志与 Token 使用。 -
结果合并:
interpreter.evaluate(...)返回ConsoleResult,其中包含此次执行的ChainResult。-
DCChains 调用
chain_res += console_res,触发ChainResult.__iadd__:- 合并
content、origin、current_generate等字段; - 将所有子链的
usage通过merge_usages汇总; - 把生成的
intermediate_msgs、tool_returns附加到共享上下文。
- 合并
-
成功返回或最终失败:
- 任意一次执行成功即返回累计后的
ChainResult; - 超过
max_plan_times仍失败则抛出ValueError,并保留失败脚本供排查。
__chains_intermediate_results 中间结果传递机制¶
设计背景¶
tf_plan 脚本在一次执行中往往需要串联多个子链,后序子链需要读取前序子链的输出来决定自身的行为。但每个子链在调用时接收的是一个独立构建的 current_input 对象(见运行流程中的"脚本解析执行"),不会自动继承上一个子链的输出。
__chains_intermediate_results 机制正是为此而设计:它以一个特殊 key 将所有子链的结构化输出累积存入同一个 dict,并通过精确的同步协议在整个 TFChainInterpreter 生命周期内保持可见。这样,后序子链可以通过其 Prompt 模板直接引用前序子链产出的任意字段,无需 tf_plan 脚本显式传递变量。
写入:Chain 进入 succeeded 时¶
每个子链执行成功后,在其状态机的 on_enter_succeeded 回调中自动写入,写入条件由 catch_intermediate_chain_result 开关控制(默认 True)。
写入逻辑是追加合并而非覆盖:
# chain.py — on_enter_succeeded() 核心逻辑
chains_results = json.loads(current_input.additional_kwargs.get("__chains_intermediate_results", "{}"))
chains_results.update(this_chain_json_output) # 追加,不替换
current_input.additional_kwargs["__chains_intermediate_results"] = json.dumps(chains_results)
写入内容取决于子链的输出类型:
| 子链输出类型 | 存储方式 | 说明 |
|---|---|---|
有 output_json_schema 的结构化输出 |
LLM 返回的 dict 直接合并 | 字段名由 schema 决定 |
纯字符串输出(设置了 result_cache_key) |
{result_cache_key: "字符串内容"} |
包装为单字段 dict |
纯字符串且无 result_cache_key |
不写入 | 该子链结果不进入缓存 |
传递路径:两段式同步¶
TFChainInterpreter 通过一种非对称的两段式同步协议保证缓存在子链之间传递,而不污染其他参数。
调用前——浅拷贝隔离
每次调用子链时,construct_chain_run_input() 对 self.additional_info 做浅拷贝,再将本次调用的 kwargs 合并进去,最后用拷贝构造 current_input:
additional_info_copy = self.additional_info.copy() # 浅拷贝
additional_info_copy.update(this_call_kwargs) # 合并本次调用参数
current_input = TextMessage(additional_kwargs=additional_info_copy, ...)
拷贝的目的是隔离:子链 A 的 param1="xxx" 不应出现在子链 B 的 additional_kwargs 里。如果直接共享引用,任何一个子链写入的参数都会污染后续所有子链。
调用后——手动同步 __chains_intermediate_results
子链执行完成后,visit_FunctionCallNode() 检查 current_input.additional_kwargs 中是否有新的缓存内容,如果有则手动同步回 self.additional_info:
if CHAINS_INTERMEDIATE_RESULTS in current_input.additional_kwargs:
self.additional_info[CHAINS_INTERMEDIATE_RESULTS] = \
current_input.additional_kwargs[CHAINS_INTERMEDIATE_RESULTS]
其他 kwargs 不同步——这是有意为之。__chains_intermediate_results 是唯一被特殊对待的 key,拥有"隔离写入、显式同步"的语义,而其他参数始终保持本次调用的独立性。
完整时序
TFChainInterpreter 初始化
└── self.additional_info = brain_ctx.current_input.additional_kwargs(引用传递)
调用子链 A:
1. additional_info_copy = self.additional_info.copy()
2. chain_A.run(current_input=TextMessage(additional_kwargs=additional_info_copy))
3. chain_A.on_enter_succeeded():
additional_info_copy["__chains_intermediate_results"] = '{"field_a": ...}'
4. interpreter 手动同步:
self.additional_info["__chains_intermediate_results"] = '{"field_a": ...}'
调用子链 B:
1. additional_info_copy = self.additional_info.copy()
→ 此次拷贝已含 {"__chains_intermediate_results": '{"field_a": ...}'}
2. chain_B.run(current_input=TextMessage(additional_kwargs=additional_info_copy))
3. chain_B 内部 Prompt 可通过 ChainsIntermediatePrompt 读取 field_a
4. chain_B.on_enter_succeeded():
additional_info_copy["__chains_intermediate_results"] = '{"field_a": ..., "field_b": ...}'
5. interpreter 手动同步:
self.additional_info["__chains_intermediate_results"] = '{"field_a": ..., "field_b": ...}'
Chain 内部读取:ChainsIntermediatePrompt¶
子链通过专用 Prompt 类 ChainsIntermediatePrompt(tfrobot.brain.chain.prompt.chain_intermediate_prompt)读取缓存内容。该类从 current_input.additional_kwargs["__chains_intermediate_results"] 中反序列化出 dict,以 chains_intermediate_result(单数)为模板变量名注入 Prompt 模板:
{# Jinja2 模板示例 #}
上一步的分析结果:{{ chains_intermediate_result.summary }}
请基于以上内容,进一步提取关键指标。
# F-String 模板示例
"上一步的搜索结果:{chains_intermediate_result.search_result}"
ChainsIntermediatePrompt 有内置的参数校验:如果 dict 中缺少模板所需的字段,整个 Prompt 段渲染为空字符串而不是抛出异常,避免因缓存不完整而导致子链崩溃。
生命周期边界¶
| 场景 | 行为 |
|---|---|
同一 tf_plan 脚本内多个子链串行 |
累积追加,子链 N 可读取子链 1…N-1 的所有结果 |
| SeqChains 容器内的子链 | 共享同一 current_input 对象,缓存自动传递,无需 interpreter 介入 |
| ParallelChains 容器内的子链 | 同样共享 current_input 对象,缓存传递有效,但并发写入时存在竞态风险 ⚠️ |
| 同一 DCChains 嵌套多层 | 各层创建独立 TFChainInterpreter,但 additional_info 来自同一 current_input.additional_kwargs 引用,缓存跨层共享 |
跨 tf_plan 脚本(多轮 exec_plan) |
brain_ctx.current_input 是同一对象,前轮写入的缓存在后轮中仍然可见,可用于构建渐进式求解;但也意味着旧结果不会自动清理 |
| 脚本执行失败后重新规划 | 失败前已写入的缓存不回滚,下次重新规划时仍可读取,有助于复用已完成的中间步骤 |
使用建议¶
- 命名要唯一:同一 DCChains 中的多个子链若都写入缓存,各自的字段名(由
output_json_schema的 property 名或result_cache_key决定)必须不同,否则后写入的会覆盖先写入的。 - 避免并行写入冲突:在
ParallelChains中,如果多个子链同时写入__chains_intermediate_results,最终结果取决于写入顺序,行为不确定。对有依赖关系的子链,使用SeqChains保证顺序。 - 不依赖跨轮缓存的持久性:跨
tf_plan轮次的缓存传递是副作用而非设计契约。如果 plan 需要隔离执行,应在构造新current_input时主动清除该 key。 - 字符串输出必须设置
result_cache_key:没有结构化 schema 的子链(如纯文本摘要链)不会自动写入缓存,需显式设置result_cache_key才能让后续子链读取其结果。
usage 统计与合并¶
- Interpreter 层聚合:
-
TFChainInterpreter在每次子链运行后把ChainResult.usage存入内部列表,并在evaluate/aevaluate收尾阶段调用merge_usages汇总 Token 消耗。 -
DCChains 层累积:
-
DCChains 在
chain_res += console_res时会将 interpreter 返回的token_usage继续累加到顶层ChainResult.usage。因此调用方获取到的result.usage已经包含本轮脚本中所有子链与 LLM 调用的整体预算。 -
持续失败的记录:
- 即便脚本失败并被加入
failed_plans,当轮执行的 Token 消耗也已被记录,方便后续分析问题或优化规划逻辑。
使用建议¶
- 脚本安全性:限制
orchestrate_chain输出的 DSL 仅包含允许的函数/链名称,避免执行未授权的行为。 - 缓存命名约定:若多条链写入
_CHAINS_INTERMEDIATE_RESULTS,请为result_cache_key设计清晰的命名模式,防止键名冲突。 - 监控 usage:复杂脚本可能在规划失败后持续重试,建议结合
max_plan_times与 Token 用量监控,必要时引入提前中断策略。 - 复用工具集:脚本可调用
tools列表中的工具。对于 DSL 中需要重复使用的能力,建议包裹为 Chain,再由 DCChains 统一调度,便于结果缓存与统计。 - 调试日志:
TFChainInterpreter会把执行过程写入logs,在ChainResult中可查看完整轨迹,对排错非常重要。
异步执行说明¶
_async_run与同步版本逻辑一致,差异仅在使用interpreter.aevaluate与异步链的async_run。- 同样支持
_CHAINS_INTERMEDIATE_RESULTS与usage聚合,适用于需要并发调用子链或工具的场景。