Skip to content

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 回调

  1. 入参包含:当前链上下文、DCChains 实例、历史失败的 tf_plan 列表与错误信息。实现可以基于这些信息决定是否重新规划或降级策略。
  2. 返回值必须是形如 tf_plan\n...\n 的字符串,运行前 DCChains 会去除外层标记并交给 Lark + TFASTTransformer 解析。
  3. 默认实现 default_orchestrate_chain 使用 DeepSeek 模型生成脚本,自动把可用子链的输入/输出 schema 注入到模板中。

运行流程

  1. 上下文构造
  2. _run / _async_run 首先把传入的 current_inputconversationknowledgetoolsintermediate_msgs 封装到 ChainContext
  3. 初始化空的 ChainResult,用于累计执行结果与 Token 统计。

  4. 规划与重试

  5. 最多 max_plan_times 次循环调用 orchestrate_chain 获取最新脚本。
  6. 每次失败都会把原脚本与失败原因写入 failed_plans / failed_infos,供下一轮规划参考。

  7. 脚本解析执行

  8. 使用 Lark(PYTHON_GRAMMAR) 解析脚本,TFASTTransformer 转换为 AST。
  9. TFChainInterpreter 负责遍历 AST:根据脚本调用 chains_map 内注册的子链、执行工具、记录日志与 Token 使用。

  10. 结果合并

  11. interpreter.evaluate(...) 返回 ConsoleResult,其中包含此次执行的 ChainResult
  12. DCChains 调用 chain_res += console_res,触发 ChainResult.__iadd__

    • 合并 contentorigincurrent_generate 等字段;
    • 将所有子链的 usage 通过 merge_usages 汇总;
    • 把生成的 intermediate_msgstool_returns 附加到共享上下文。
  13. 成功返回或最终失败

  14. 任意一次执行成功即返回累计后的 ChainResult
  15. 超过 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 类 ChainsIntermediatePrompttfrobot.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 是同一对象,前轮写入的缓存在后轮中仍然可见,可用于构建渐进式求解;但也意味着旧结果不会自动清理
脚本执行失败后重新规划 失败前已写入的缓存不回滚,下次重新规划时仍可读取,有助于复用已完成的中间步骤

使用建议

  1. 命名要唯一:同一 DCChains 中的多个子链若都写入缓存,各自的字段名(由 output_json_schema 的 property 名或 result_cache_key 决定)必须不同,否则后写入的会覆盖先写入的。
  2. 避免并行写入冲突:在 ParallelChains 中,如果多个子链同时写入 __chains_intermediate_results,最终结果取决于写入顺序,行为不确定。对有依赖关系的子链,使用 SeqChains 保证顺序。
  3. 不依赖跨轮缓存的持久性:跨 tf_plan 轮次的缓存传递是副作用而非设计契约。如果 plan 需要隔离执行,应在构造新 current_input 时主动清除该 key。
  4. 字符串输出必须设置 result_cache_key:没有结构化 schema 的子链(如纯文本摘要链)不会自动写入缓存,需显式设置 result_cache_key 才能让后续子链读取其结果。

usage 统计与合并

  1. Interpreter 层聚合
  2. TFChainInterpreter 在每次子链运行后把 ChainResult.usage 存入内部列表,并在 evaluate/aevaluate 收尾阶段调用 merge_usages 汇总 Token 消耗。

  3. DCChains 层累积

  4. DCChains 在 chain_res += console_res 时会将 interpreter 返回的 token_usage 继续累加到顶层 ChainResult.usage。因此调用方获取到的 result.usage 已经包含本轮脚本中所有子链与 LLM 调用的整体预算。

  5. 持续失败的记录

  6. 即便脚本失败并被加入 failed_plans,当轮执行的 Token 消耗也已被记录,方便后续分析问题或优化规划逻辑。

使用建议

  1. 脚本安全性:限制 orchestrate_chain 输出的 DSL 仅包含允许的函数/链名称,避免执行未授权的行为。
  2. 缓存命名约定:若多条链写入 _CHAINS_INTERMEDIATE_RESULTS,请为 result_cache_key 设计清晰的命名模式,防止键名冲突。
  3. 监控 usage:复杂脚本可能在规划失败后持续重试,建议结合 max_plan_times 与 Token 用量监控,必要时引入提前中断策略。
  4. 复用工具集:脚本可调用 tools列表中的工具。对于 DSL 中需要重复使用的能力,建议包裹为 Chain,再由 DCChains 统一调度,便于结果缓存与统计。
  5. 调试日志TFChainInterpreter 会把执行过程写入 logs,在 ChainResult 中可查看完整轨迹,对排错非常重要。

异步执行说明

  • _async_run 与同步版本逻辑一致,差异仅在使用 interpreter.aevaluate 与异步链的 async_run
  • 同样支持 _CHAINS_INTERMEDIATE_RESULTSusage 聚合,适用于需要并发调用子链或工具的场景。