Vibe-Trading 仓库的多智能体运行时:任务怎么派,崩了怎么办

2026-08-05

本文基于 Vibe-Trading 仓库 commit 3a752d5(2026-08-04)梳理,该项目仍在高频迭代,具体行为以仓库 https://github.com/HKUDS/Vibe-Trading 最新代码与文档为准。

很多人读多智能体项目只读编排那一个文件,然后得出”不过就是个拓扑排序”的结论——在 Vibe-Trading 这个开源项目里,真正决定你半夜会不会被叫起来的,是编排之外那三块:worker 怎么判断自己算不算干完了、任务状态落到哪个文件才算数、以及读接口把什么信息藏起来了。 这四块合起来才是一套可用的运行时,缺任何一块,另外三块的设计都会显得莫名其妙。

先做个消歧:Vibe-Trading 是 HKUDS 放出来的开源仓库名,不是”凭感觉交易”那种说法。仓库许可证是 MIT(Copyright 2026 Vibe-Trading Contributors)。本文只讨论它的 Agent 工程实现,不讨论任何交易决策。

站内相近的几篇分工不同:多 Agent 框架选型是横向比几套框架该选哪个,Agent 并发编排讲的是不绑定具体项目的并发编排通用原则,Superpowers 的并发派遣拆的是另一套技能体系里的派单机制;本篇只钻 Vibe-Trading 一个仓库的 agent/src/swarm/ 目录,把代码逐行读出来的结论写下来。

一、四块各管什么

先把地图摆平。这四个文件不是平级模块,是一条链路上的四段。

组成部分它负责什么仓库位置你什么时候会碰到它
SwarmRuntime校验 DAG、算执行分层、起后台线程、逐层调度、汇总收尾agent/src/swarm/runtime.py改并发度、改重试策略、查”为什么这个任务没跑”
run_worker单个 worker 的 ReAct 循环、提示词拼装、工具执行、交付物判定agent/src/swarm/worker.py查”这个 agent 为什么被判不合格”、加工具、调迭代预算
TaskStore 与 DAG 算法每个任务一个 JSON 文件的读写、依赖解锁、环检测、分层agent/src/swarm/task_store.py排查任务状态与预期不符、写外部监控
serialize_task把任务投影成对外读接口的字典,统一字段口径agent/src/swarm/serialization.py对接 API、写前端、发现”失败了但看不到原因”

再补两个上下文。运行目录结构写在 agent/src/swarm/store.py 的模块注释里:.swarm/runs/{run_id}/ 下面是 run.jsonevents.jsonltasks/inboxes/artifacts/。编制文件放在 agent/src/swarm/presets/,我数了一下是 30 份 yaml(ls 一遍就能复现),文件名从 equity_research_team.yamlrisk_committee.yaml 都有,一份 yaml 就是一支团队的角色表加任务图。

二、任务怎么派:先验环,再分层,再起线程

start_run 这个入口做的事比名字多。它先调 self._store.reap_stale_running_runs() 打扫上一轮的残局——把那些宿主进程死掉、状态却还挂在 running 的历史 run 收掉;这一步包在 try 里,扫墓失败只记 warning,不影响本次启动。然后 build_run_from_preset(preset_name, user_vars) 把 yaml 变成 SwarmRun,紧接着 validate_dag(run.tasks)

validate_dagtask_store.py 里是标准三色 DFS:先检查所有 depends_on 引用的任务是否真实存在,不存在直接抛 Task 'x' depends on unknown task 'y';然后染色遍历,遇到 GRAY 节点就把 path 从环起点切出来,报错信息里带完整环路径。这个细节值得学——报错只写一句”检测到环”的话,排查一个几十节点的 DAG 会很痛苦。

分层用 topological_layers,Kahn 算法,返回的是 list[list[str]]:同一层内的任务 ID 之间没有依赖关系,可以并行;层与层之间串行。这里还留了第二道环检测——processed != len(tasks) 时报 processed X/Y tasks

真正的执行体 _execute_run 跑在后台 daemon 线程里,线程名是 swarm-{run.id},所以 start_run 是立刻返回的,返回的 run 状态还是 pending。执行前它会调 invalidate_mcp_specs_cache() 强制刷新 MCP 工具发现,理由写在注释里:操作者可能在两次 run 之间改了 mcp_servers 配置,缓存的 specs 会过期。

还有一步容易漏:_prefetch_grounding_data。它从 user_vars 里抽出标的代码,预取行情数据挂到 run.grounding_data,再由 grounding.format_grounding_block 渲染成一段文本,一次性发给这次 run 的每个 worker。抽取数量有上限,超了就截断并打 warning,上限来自 SWARM_GROUNDING_MAX_SYMBOLS。这一步可能耗时几十秒,所以整段被 HeartbeatTimer 包住,持续往 events.jsonl 里写 run_heartbeat 事件——不这么做,后面要讲的僵尸 run 回收器会把一个正在正常拉数据的健康 run 误杀。

层内并行靠 ThreadPoolExecutor(max_workers=self._max_workers)。注意它没有用 with,而是手动管生命周期,finally 里调 executor.shutdown(wait=False, cancel_futures=True)。注释说得很直白:用 with 的话 KeyboardInterrupt 会卡在 shutdown(wait=True) 上,CLI 按了 Ctrl-C 半天退不出去。

三、worker 那一圈:干完了不等于交付了

run_worker 没有复用项目自己的 AgentLoop,而是直接拿 ChatLLM 加一个 for 循环手搓 ReAct,文件头注释写明这是为了”保持 worker 自包含、agent 核心不动”。

循环体里每轮固定做几件事。开头先做微压缩:把消息里的 tool 结果,除最近 _KEEP_RECENT_TOOLS 条以外、且长度超过阈值的,内容直接换成 [cleared]。然后检查两条硬边界——经过时间超过 timeouttimeout 分支,序列化后的粗估 token 超过 _MAX_TOKEN_ESTIMATEtoken_limit 分支。到迭代预算 80% 的位置,会往消息里塞一条 [SYSTEM] 开头的收尾提醒;最后一轮干脆不传 tool_defs,逼模型只能出文本。

模型调用统一走 llm.stream_chat,注释说非流式在某些供应商上不可靠,而且流式还能喂前端实时进度。流式失败时看 ProviderStreamError.retryable:可重试的睡一下(间隔来自 SWARM_STREAM_RETRY_DELAY_S)再试一次,确定性的 4xx 直接抛。整个 LLM 调用和每次工具执行都各自被 HeartbeatTimer 包住,发 task_heartbeat 事件并带上 phase 字段区分 llm 还是 tool

提示词由 build_worker_prompt 拼,顺序是:Role、角色自己的 system_prompt(里面的 {upstream_context} 会被上游摘要替换)、可用技能列表、grounding 块、行情工具政策(仅当该 agent 的工具白名单里有 get_market_data)、一段无条件的 Data Citation Discipline、执行规则、当前 UTC 时间。那段 Data Citation Discipline 值得单独看:它要求 worker 输出的每个具体数字都必须能追溯到本次 run 的工具结果、grounding 块或上游上下文之一,明确禁止从训练数据里回忆数字,并且特意说明这条规则同样适用于没有数据工具的汇总/编辑角色。这是把”防编造”写成了硬提示词约束,而不是指望模型自觉。

最有意思的是终态判定。WorkerStatus 有五个值:completedfailedtimeouttoken_limitincompleteincomplete 是单独立出来的一档,models.py 的注释解释了原因:worker 没抛异常地跑完了,但没产出实质交付物,绝不能和 completed混为一谈。判定逻辑在 _classify_deliverable,命中任一条就算不合格:空交付物;出现未被解析的工具调用标记(供应商没解析 tool call,标记原样吐进了文本);文本里出现 mock data 之类的自认伪造措辞;整段其实是个原始工具返回信封而不是分析;以 Phase 1 / Plan 开头且正文过短或以交接语收尾的”只有计划没有执行”的残桩;以及——只对数据类 agent 生效——一次有效工具调用都没有且没写出 report.md

最后一条为什么要限定数据类 agent,注释里也写了:像 equity_research_team 里的编辑角色,工具白名单是 bashread_filewrite_fileload_skilledit_file 这个通用集合的子集,它本来就该只产文本,按”无工具证据”判失败会误杀。判据是 _is_data_agent:工具集合减去通用集合还有剩余才算数据 agent。

四、结果怎么收:三层落盘,各有各的真源

结果不是收到内存里就完了,落盘分三层。

任务态在 tasks/task-{id}.json TaskStore.save_task 先写 .tmpreplace,加线程锁,一个任务一个文件,避免多线程写同一个大 JSON。update_status 是”读出来、改字段、验回去、再写下去”,用 Pydantic 做一次完整校验。

产物在 artifacts/{agent_id}/ worker 结束时写 summary.md,超时和迭代耗尽的路径还会额外把完整消息列表落成 messages.json 供事后复盘。摘要取值走 _resolve_summary:如果 report.md 存在且非空,就用它的内容当摘要,否则退回模型的最后一段文本。_collect_artifacts 收集时跳过软链接,并用 is_relative_to 校验解析后的真实路径没跑出 agent 产物目录,越界的一律不收。

run 级快照在 run.json 关键设计是 _sync_run_tasks_snapshot 只在层边界调用一次,注释解释是为了避免每个任务写一次带来的 I/O 噪音,并明说”每任务文件仍然是活的真源”。收尾时 run.tasks = task_store.load_all() 做最终同步,最终报告从最后一层里第一个有摘要的任务取。

对外读接口统一收敛到 serialization.pyserialize_task 返回 id、agent_id、status、summary、iterations、error、started_at、completed_at、depends_on、blocked_by。这个模块的文档字符串记录了它诞生的原因:在它之前,三个读边界各自维护字段白名单,而且三个都漏了 error——一个配错供应商的 run 在磁盘上明明记着错误原因,调用方却什么也看不到。现在三个读边界——agent/src/api/swarm_routes.pyagent/src/tools/swarm_tool.pyagent/mcp_server.py——都从这里取。run_level_error 再把第一个带错误的任务拼成 任务id/agent_id: 原因 放到顶层,让只读顶层的调用方也有信号。所有错误文本出门前都过 redact_internal_paths,不把宿主机路径泄给调用方。

五、崩了怎么办:五条兜底路径

按触发顺序排一遍。

第一条,上游没成功就不派活。 _execute_layer 在提交任务前逐个加载 depends_on 的上游任务,只要有一个不是 completed(包括文件都找不到的情况),就把当前任务标成 blocked、写明 blocked_by、发 task_blocked 事件,直接跳过不派。注释里给的例子很实在:某个编制里的组合管理角色依赖风险任务,如果风险任务失败了却照常派活,下游会拿着空上下文产出一份”结论”。同层里没有共享上游的兄弟任务不受影响。

第二条,单任务重试。 _run_worker_with_retriesagent_spec.max_retries 循环,每次重试前发 task_retry 事件并带上一次的错误。这里有个必须记住的判断:if result.status != "failed": return——只有 failed 触发重试,timeouttoken_limitincomplete 都直接返回,不再试。token 计数跨所有尝试累加,最后用 model_copy 覆写回去。

第三条,层级硬截止。 每个任务的预算是 timeout_seconds × (max_retries + 1),取全层最大值再加一个固定缓冲,作为 as_completed 的 timeout。超时的话,把还没出结果的 future 逐个 cancel(),并合成 status="timeout" 的结果。注释说明这是为了防住卡在 C 扩展或阻塞 I/O 里、绕过了循环内超时检查的 worker 线程。

第四条,取消。 cancel_run 只是 set() 一个 threading.Event,检查点在层与层之间,不是随时生效。触发后 _cancel_remaining_tasks 把所有既非 completed 也非 failed 的任务标成 cancelled。另外 KeyboardInterrupt 会顺手把同一个 event 设上再抛。

第五条,僵尸 run 回收。 store.py 里的 compute_stale_threshold 按 run 算一个”事件静默预算”:心跳间隔的十倍作为基准,上限卡在最大单任务重试预算加一分钟,下限六十秒。is_run_stale 是只读判断,reap_stale_running_runs 遍历目录逐个 reconcile。所以那些遍布代码的 HeartbeatTimer 不是装饰——它们是这套判活机制的信号源。

六、边界与代价:它明确不管什么

这套设计的取舍很清楚,用之前得认。

层间串行的代价是长尾拖累。 调度粒度是”层”而不是”任务”:某一层里九个任务两秒跑完、第十个跑满超时,整层都得等它。没有”上游一完成就立刻放行下游单点”的细粒度事件驱动。

没有断点续跑。 取消或宿主进程死掉之后,剩余任务被标 cancelled 或被回收器标掉,重跑就是从头开始。磁盘上的 tasks/*.json 是给你排查用的,不是给引擎恢复用的检查点。

单进程线程模型。 并发上限是一个 ThreadPoolExecutormax_workers,它在构造 SwarmRuntime 时传入,工具入口 agent/src/tools/swarm_tool.py 那条路径读的是 SWARM_MAX_WORKERS;没有跨机分发、没有队列中间件。任务规模上去了要自己在外面套一层。

上游上下文只有文本摘要。 input_from 映射的是 {上下文键: 上游任务ID},取的是上游的 summary 字符串,拼进下游的系统提示词。结构化数据要传,只能靠双方约定在文本里写清楚,引擎不做校验。

微压缩是有损的。 老的工具结果直接被替换成 [cleared],模型如果在第十轮想回看第二轮的原始数据,看到的是被清空的占位。

交付物判定是启发式字符串匹配。 前缀表、伪造词表、交接语尾表都硬编码在 worker.py 里。换个语种或换个写作风格的模型输出,误判方向两头都有可能。

它明确不管的事: 模型输出的分析质量本身;下游任何形式的执行安全。仓库里 agent/src/trading/connectors/ 下有 12 个券商连接器子目录(README 也自述 12 brokers),一旦链路接到真实券商,凭据的暴露面、下错单不可撤销、程序化交易在不同司法辖区的合规义务差异,都不是这套 swarm 运行时能替你兜的——这几件事以你所在司法辖区的监管要求与券商协议为准。

顺带说清一件常被误读的事:仓库根目录 NOTICE 明确写了几个因子库各自的上游来源与许可——Microsoft Qlib 的特征定义走 Apache 2.0,另有几组公式来自公开论文与研报,仓库把它们当作数学事实重实现,各因子库子目录下另有 LICENSE.md。所以它们不是”项目自研的因子库”,能不能商用以许可证原文为准,本文不提供法律意见。这些因子的历史表现不代表未来,本文只讨论工程实现。

七、上手与避坑清单

别指望 incomplete 会被自动重试。 为什么会踩:直觉上”没交付”就是”失败”,你以为配了 max_retries 就有两次机会。实际上重试只认 failedincomplete 一次就终结,而且它会让下游被 blocked。怎么避:把 incomplete 当成配置问题而不是抖动问题去查——它通常意味着提示词没让 agent 写出 report.md,或者工具白名单给了数据 agent 却没给它能真正调通的数据工具。

别把 run.json 当实时状态源。 为什么会踩:你写了个轮询脚本读 run.json,发现某个任务显示 in_progress 好几分钟没动。实际上那个文件只在层边界刷新,代码注释直接写了”每任务文件仍然是活的真源”。怎么避:要实时就读 tasks/task-*.json,或者订阅 events.jsonl

update_status 传参数时把键名核对一遍。 为什么会踩:它的实现是 for key, value in kwargs.items(): if key in updated_data——键名不在模型字段里就被静默丢弃,不报错、不告警。怎么避:改这块代码时对着 models.pySwarmTask 字段表核一遍,或者补一条单测断言写进去的值能读回来。

留意产物目录是按 agent_id 分的,不是按 task_id。 为什么会踩:artifact_dir = run_dir / "artifacts" / agent_id。同一个 agent 在一次 run 里承担多个任务时,它们共用同一个目录,而摘要解析又优先读 report.md。怎么避:编制里让一个 agent 只担一个任务;确实要复用角色,就在提示词里要求写到不同文件名,或者在派活前确认上一轮产物已经归档。

include_shell_tools 打开之前先想清楚。 为什么会踩:这个开关从 start_run 一路透传到 build_swarm_registry,默认是 False;打开就等于给每个 worker 一个可执行 shell,而 worker 的行为由模型输出决定。怎么避:只在隔离环境里开;事件里的参数预览虽然过了 is_sensitive_argredact_payload 做脱敏,但那只保护日志,不保护进程内真实凭据。

模板变量缺了不会报错。 为什么会踩:_FallbackDict.__missing__ 会把缺失的键换成一句”请根据目标自行判断”的英文提示塞进提示词,看起来一切正常,实则模型在自由发挥。怎么避:调 run_swarm 前把编制里用到的占位符和你传的 user_vars 对一遍键名。

心跳配置别乱调。 为什么会踩:SWARM_HEARTBEAT_INTERVAL_S 同时被 worker 的心跳器和 compute_stale_threshold 读取,两边必须一致;把它调得很大,僵尸 run 的检测延迟也跟着变大。怎么避:认准它是判活参数不是日志频率参数,要改就连着回收阈值一起想。

收尾:接下来该读哪里

一条自检路线,照着走一遍就能把这套运行时的行为吃透:拿 agent/src/swarm/presets/ 里任意一份 yaml,画出它的 depends_on 图,手算一遍 topological_layers 的分层结果;然后跑一次,对着 events.jsonl 核事件序列是不是 run_startedlayer_startedtask_startedtool_call/tool_resulttask_completed;再故意把某个中间任务的工具白名单改错,看下游是不是如期变成 task_blocked 而不是拿着空上下文硬跑。

读完这四个文件还想往下钻,agent/src/swarm/store.py 是下一站——原子写、Windows 上 os.replace 的重试退避、reconcile_run 的对账逻辑都在那里,也是理解”崩了怎么办”的最后一块。想看这套东西怎么暴露给调用方,就跟着 serialize_task 的三个引用点走到 agent/src/api/swarm_routes.pyagent/src/tools/swarm_tool.pyagent/mcp_server.py

重试与失败分类这块的通用套路,可以对照失败重试的设计可观察日志怎么设计一起看——Vibe-Trading 的做法是其中一种取舍,不是唯一解。

本篇属于一个把开源个人交易 Agent 项目 Vibe-Trading逐层拆开讲的系列,整体地图见 Vibe-Trading 是什么:HKUDS 这个开源交易 Agent 项目的工程全景与边界;沿着这条线往下,还可以看 Vibe-Trading 假设注册表:开源交易 Agent 怎样拦住越研究越自信开源项目 Vibe-Trading:30 份 yaml 定义的多智能体团队编制

想系统学会用 AI?报名体系课或加入会员,照着学、照着用。