数据截至 (上游 commit 7a975c596eca)
运行时与事件流 — 点一下 Run 之后前后端之间发生了什么
30 秒导读: 在 Langflow 画布上点 Run,浏览器并不是发一个请求等一个结果。它先发一个 POST 拿到
job_id,再开一条长连接把后端吐出来的事件流一条条读回来;后端一边跑图一边把「排序完成 / 某个点开始 / 某个点结束 / 又来一个 token / 结束」写进一个队列,由这条连接推给前端。本章把这条链路从按钮一路追到asyncio.Queue,再追回节点变色。
本章不讲组件怎么被发现、怎么变成 MCP 工具(那是 06),也不细讲"谁先跑、谁能并行"的调度规则(那是 03)。这里只关心一次交互式运行的端到端链路和事件协议。
1. 先看全貌:一次 Run 的三段旅程
一次交互式运行被切成三个独立的 HTTP 交互,而不是一个。
浏览器 后端
│ │
│ ① POST /api/v1/build/{flow_id}/flow│
├───────────────────────────────────>│ 建 job_id + asyncio.Queue
│ │ 起后台 task 跑 generate_flow_events
│ <──────── {"job_id": "..."} ───────┤ │
│ │ │ 边跑边把事件塞进队列
│ ② GET /api/v1/build/{job_id}/events│ ▼
├───────────────────────────────────>│ ┌───────────────┐
│ <═══ NDJSON 事件流(长连接) ═════════┤ │ asyncio.Queue │
│ vertices_sorted / build_start │ └───────────────┘
│ end_vertex / token / end │
│ │
│ ③ POST /api/v1/build/{job_id}/cancel(用户点停止时)
├───────────────────────────────────>│ event_task.cancel()
三段各自的职责:
| 阶段 | 路由 | 干什么 | 定义在哪 |
|---|---|---|---|
| ① 建 job | POST /build/{flow_id}/flow | 建队列、起后台任务、立刻返回 job_id,不等结果 | src/backend/base/langflow/api/v1/chat.py:256 (build_flow) |
| ② 拉事件 | GET /build/{job_id}/events | 把队列里的事件按投递方式送给浏览器 | src/backend/base/langflow/api/v1/chat.py:414 (get_build_events) |
| ③ 取消 | POST /build/{job_id}/cancel | 取消那个后台任务 | src/backend/base/langflow/api/v1/chat.py:441 (cancel_build) |
三个路由的实现全部落在同一个文件里:src/backend/base/langflow/api/build.py 的 start_flow_build、get_flow_events_response、cancel_flow_build。
为什么要拆成两个请求? 因为「启动」和「读结果」的失败语义完全不同。POST 失败意味着这张流根本没跑起来(校验不过、组 件被禁用),要立刻报错;GET 失败只意味着这条连接断了——任务还在后端跑,重连就能接着读。拆开之后,重连不需要重跑。
2. 第一段:POST 只做一件事——建 job
start_flow_build 短得出奇,因为它真正只做三件事。
job_id = str(uuid.uuid4())
_, event_manager = queue_service.create_queue(job_id) # 1. 建队列 + 事件管理器
task_coro = generate_flow_events(...) # 2. 准备协程(还没跑)
queue_service.start_job(job_id, task_coro) # 3. 丢进后台
return job_id
— src/backend/base/langflow/api/build.py:182-250,符号 start_flow_build。
三件 事对应到 JobQueueService 的两个方法:
create_queue(src/backend/base/langflow/services/job_queue/service.py:178)建一个无界asyncio.Queue,并用它构造一个EventManager,登记进self._queues[job_id]。start_job(同文件:206)用asyncio.create_task把协程跑起来,任务句柄存回同一条记录。
start_job 有个不显眼但关键的包装:协程被 _guarded_task 裹了一层(:246)。它的作用是——万一 generate_flow_events 抛了没接住的异常,也要保证队列里先落一条 error 事件、再落一个结束哨兵 (None, None, ts) 才退出。没有这层包装,消费端会永远卡在 queue.get() 上等一个永远不来的结束信号。
顺带澄清一个容易混淆的文件名:
src/backend/base/langflow/api/v1/flow_events.py不是本章说的作业队列。那是另一套东西——GET/POST /api/v1/flows/{flow_id}/events,给某张流追加/读取"设计期通知事件"(FlowEventsService,src/backend/base/langflow/services/flow_events/service.py:44),和运行时的构建事件流没有关系。运行时的作业队列在services/job_queue/service.py。
3. 第二段:同一批事件,三种投递方式
get_flow_events_response(build.py:253)根据查询参数 event_delivery 分岔。三种模式的取舍:
| 模式 | 值 | 传输形态 | 谁在用 | 代价 |
|---|---|---|---|---|
| 流式 | streaming | 一次 GET,长连接,NDJSON 边算边推 | 前端默认 | 需要中间层不缓冲响应 |
| 轮询 | polling | 反复 GET,每次把队列排空返回,25ms 后再来 | 流式不可用时的降级 | 延迟 = 轮询间隔 |
| 直连 | direct | POST /build 本身就返回流,不给 job_id | 前端 e2e 测试 | 无法重连 |
枚举定义在 src/backend/base/langflow/api/utils/core.py:85(EventDeliveryType),前端同名枚举在 src/frontend/src/constants/enums.ts:42。
┌── streaming ─ 一次 GET,长连接,事件产生即推
生产端 asyncio.Queue ────┼── polling ─── 反复 GET,每次 drain 光队列
└── direct ──── POST 直接变成流
3.1 流式:create_flow_response 的三个零件
create_flow_response(build.py:342)返回一个 DisconnectHandlerStreamingResponse,内部有三个协作的零件。
零件一,消费循环 consume_and_yield(build.py:392-407)。就是 await queue.get() 然后 yield _project_event_to_v1(value.decode("utf-8")),直到取到 value is None 这个哨兵才 break。事件在队列里放了多久也顺手记进 debug 日志。
零件二,心跳 _heartbeat(build.py:367-382)。每 STREAMING_ACTIVITY_REFRESH_S(默认 10 秒,build.py:60)调一次 queue_service.touch_activity(job_id)。注释里点明了它为什么必须独立于事件节奏:一次慢 LLM 调用可能几十秒不产事件,如果用"最近一次事件时间"来判活,看门狗会把一个正常跑着的 job 当成僵尸回收。心跳把"客户端还在"和"任务还在产出"这两件事解耦了。
心跳只在后端支持时才起:touch = getattr(queue_service, "touch_activity", None),内存版 JobQueueService 根本没这个方法,于是 heartbeat_task 为 None,整段逻辑自动消失(build.py:365,384)。touch_activity 只有 Redis 版实现(services/job_queue/service.py:1081)。
零件三,断线处理 on_disconnect(build.py:409-430)。浏览器一关,DisconnectHandlerStreamingResponse.listen_for_disconnect(src/backend/base/langflow/api/disconnect.py)收到 http.disconnect,回调就会:停心跳 → event_task.cancel() 取消构建 → 清队列 → 发一个 on_end。用户关标签页,后端就真的停算了,不会把一次没人要的 LLM 调用跑完。
多进程部署下还有一条支线:如果这条 GET 落在了没拥有那个 task 的 worker 上(event_task is None),就退而求其次调 queue_service.signal_cancel(job_id),通过 Redis 频道把取消信号发给真正的 owner(build.py:414-424)。
3.2 轮询:一次请求排空一次队列
轮询分支(build.py:295-330)的做法很直白:
while not main_queue.empty(): # 有多少拿多少,不阻塞
_, value, _ = await main_queue.get()
if value is None: ... # 哨兵 → 取消 event_task、发 on_end
events.append(value.decode("utf-8"))
if not events: # 一条都没有 → 阻塞等第一条,避免空转
_, value, _ = await main_queue.get()
return Response(content="\n".join(...), media_type="application/x-ndjson")
注意 if not events 那一支:如果这一轮什么都没有,它会阻塞等一条再返回。这让"没有事件"的空转成本被摊薄到服务端的一次 await,而不是客户端疯狂空跑。
event_delivery 落到未知值时不会静默按轮询处理,而是直接 400 并把支持的值列出来(build.py:277-293)——注释说明了理由:三种模式的跨 worker 保障不同,静默 fallthrough 会掩盖多 worker 配置错误。
4. 第三段:后端到底在跑什么
generate_flow_events(build.py:439)是整台机器的主体。它内部定义了一串闭包,链路是:
generate_flow_events
│
├─ build_graph_and_get_order() 建图 + 排序
│ ├─ create_graph() build_graph_from_db / build_graph_from_data
│ └─ sort_vertices() graph.sort_vertices(stop_id, start_id)
│
├─ on_vertices_sorted ← 第一条事件:告诉前端「这些点要跑」
├─ on_build_start
│
└─ _run_vertex_build()
└─ 对第一层每个 id: create_task(build_vertices(id))
│
└─ _build_vertex → graph.build_vertex
然后 on_end_vertex
然后对 next_vertices_ids 递归 create_task
4.1 建图与排序
create_graph(build.py:615-659)二选一:
- 请求体没带
data→ 从数据库读那张流:build_graph_from_db(build.py:625)。 - 带了
data(画布上未保存的改动)→ 直接用请求体建:build_graph_from_data(build.py:644)。
这就是"画布上改了没保存也能点 Run"的实现。图怎么从 JSON 变出来,见 02-graph-construction。
sort_vertices(build.py:661-666)调 graph.sort_vertices(stop_component_id, start_component_id),拿到第一层可跑的点。它外面套了一层 try/except:带 start/stop 的排序失败时,退回无参数的全图排序,而不是让整次运行崩掉。
排完序立刻发第一条事件(build.py:928):
event_manager.on_vertices_sorted(data={"ids": ids, "to_run": vertices_to_run})
ids 是第一层,to_run 是整次运行会碰到的全部点——前端用后者把所有相关节点先标成 TO_BUILD(灰色待跑),用前者标成"马上开始"。
4.2 并发是"铺开"出来的,不是排好的
这是本章最值得带走的一处设计。build_vertices(build.py:833-883)结构如下:
vertex_build_response = await _build_vertex(vertex_id, graph, event_manager) # 跑这个点
...
event_manager.on_end_vertex(data={"build_data": build_data}) # 报告结果
if vertex_build_response.valid and vertex_build_response.next_vertices_ids:
tasks = [asyncio.create_task(build_vertices(nid, ...)) # 对每个后继开一个 task
for nid in vertex_build_response.next_vertices_ids]
await asyncio.gather(*tasks)
next_vertices_ids 来自 graph.get_next_runnable_vertices(lock, vertex=vertex, cache=False)(build.py:697)。所以并发度不是谁事先算好的,而是每个点跑完自己报告"接下来谁能跑",然后就地为每个后继开一个 task,并发自然铺开:
第一层 ids = [A, B]
├─ task(A) ── 完 ── next=[C] ──── task(C) ── 完 ── next=[E] ── task(E)
└─ task(B) ── 完 ── next=[C,D] ─┬ task(C) (已在跑的会被 run_manager 挡掉)
└ task(D)
"谁算可跑、重复的怎么去重、环怎么办"全部封在 graph.get_next_runnable_vertices 和 RunManager 里——见 03-scheduler 和 04-cycles-and-branching。运行时这一层只负责照单开 task。
对比一下同一个引擎的另一种驱动方式会更清楚。非交互路径用的 Graph.process(src/lfx/src/lfx/graph/graph/base.py:2187-2248)是层同步的:一层的 task 全 gather 完,_execute_tasks 返回下一层,才开下一批。两种驱动的差别:
/build(交互) | Graph.process(/run) | |
|---|---|---|
| 驱动方式 | 每点跑完就地递归 spawn | 攒成一层,gather 完再开下一层 |
| 层间屏障 | 无 | 有 |
| 后果 | 快的分支不等慢的分支 | 一层里最慢的点拖住整层 |
| 代码 | build.py:867-883 | graph/graph/base.py:2216-2248 |
4.3 单点执行 _build_vertex
_build_vertex(build.py:672-831)是"跑一个点"的完整包裹,核心一句是 await graph.build_vertex(...)(build.py:684),把 event_manager 一路传进去 (graph/graph/base.py:2049)。它额外做的事:
- 错误不炸流程:组件抛异常时不往上抛,而是把 traceback 包成
OutputValue(type="error")塞进result_data_response,valid=False(build.py:705-732)。因为这个点失败了,后继不该跑,但事件流必须继续,前端才能把这个节点涂红。 - 落库与缓存:不流式且
log_builds为真(且graph.persist_messages为真——匿名运行会把它置假、整体跳过落库)→ 后台任务log_vertex_build落库;否则把整张图写进 chat 缓存(build.py:743-757)。 - 停点裁剪:如果用户指定了停在某个组件,且它就在后继里,就把后继裁成只剩它(
build.py:777-778)。 - 计时:
duration/timedelta写回响应,前端拿它显示每个节点跑了多久。
最后收尾(build.py:1007-1009):把所有点的耗时求和,发 on_end(data={"build_duration": ...}),再往队列里放结束哨兵 (None, None, time.time())(build.py:1032)。哨兵是流结束的唯一信号,前面提到的 _guarded_task 就是为了保证异常路径上它也一定会被放进去。
5. 事件对象:EventManager 与它的三层容错
前面反复出现的 event_manager.on_xxx(...) 到底是什么?EventManager(src/lfx/src/lfx/events/event_manager.py:30)是一个把"发生了什么"翻译成队列里一行 JSON 的适配器——只有大约 120 行,但每一处容错都有故事。
5.1 事件清单
create_default_event_manager(event_manager.py:124-136)登记了全部 10 种事件。这就是前后端之间的完整协议:
| 方法名 | 线上事件名 | 谁发 | 前端拿它干什么 |
|---|---|---|---|
on_vertices_sorted | vertices_sorted | build.py:928 | 把待跑节点标成 TO_BUILD |
on_build_start | build_start | build.py:931 / 组件 | 无 id 时记录起跑时刻;有 id 时标 BUILDING |
on_end_vertex | end_vertex | build.py:867 | 主力事件:写结果、涂色、点亮下一批边 |
on_build_end | build_end | 组件 | 标 BUILT |
on_token | token | 组件流式输出 | 往聊天气泡追加字符 |
on_message | add_message | 组件 | 新增一条聊天消息 |
on_remove_message | remove_message | 组件 | 撤回一条消息 |
on_log | log | custom_component/component.py:1782 | 追加日志到 flowPool |
on_error | error | build.py:899,969 | 弹错误、标红 |
on_end | end | build.py:1008 | 收尾:算总时长、isBuilding=false |
前端 onEvent 的 switch 分支(src/frontend/src/utils/buildUtils.ts:721-849)和这张表一一对应,没有多余分支也没有遗漏。
5.2 注册:register_event 用 partial 把 event_type 焊死
def register_event(self, name, event_type, callback=None):
if not name.startswith("on_"): raise ValueError(...)
callback_ = partial(self.send_event, event_type=event_type) # 把类型名预先绑上
self.events[name] = callback_
— event_manager.py:54-70。
调用方于是只需要写 on_end_vertex(data=...),不必每次重复事件类型字符串。名字必须以 on_ 开头是硬约束。
5.3 发送:send_event 的三层容错
send_event(event_manager.py:72-115)短短 40 行里叠了三层"绝不因为发事件而搞崩流程"的保护。
第一层:序列化失败就降级,不抛。
try:
jsonable_data = jsonable_encoder(data)
except (ValueError, TypeError) as exc:
jsonable_data = serialize(data, to_str=True) # 兜底:不认识的对象转字符串
— event_manager.py:73-84。注释点名了触发场景(issue #12591):某个向量库客户端里揣着一个 threading.Lock,通过 Loop 组件的每轮构建事件冒出来,FastAPI 的 jsonable_encoder 直接抛 ValueError。如果不降级,整次流程会因为"一个日志字段没法转 JSON"而挂掉。
第二层:跨线程用 call_soon_threadsafe。
if in_event_loop:
self.queue.put_nowait(item)
elif self._loop is not None and self._loop.is_running():
self._loop.call_soon_threadsafe(self.queue.put_nowait, item)
else:
self.queue.put_nowait(item) # 无 loop 的同步环境,如单测
— event_manager.py:97-111。self._loop 在构造时抓一次(event_manager.py:35)。
这不是防御性编程的空转——真有调用方在别的线程里。token 事件就是:Component._process_chunk 用 await asyncio.to_thread(self._event_manager.on_token, ...) 把它扔到线程池里发(src/lfx/src/lfx/custom/custom_component/component.py:2122-2128)。直接 put_nowait 虽然也能进队列,但唤不醒在 queue.get() 上等着的那个协程。
第三层:队列满只丢事件,不炸流程。
except asyncio.QueueFull:
logger.warning("Event queue full; dropping event_type=%s", event_type)
— event_manager.py:112-113。取舍很明确:丢一条 token 的可观测性 ≪ 让整次运行失败。
5.4 兜底:没注册的事件是个 noop
def noop(self, *, data): pass
def __getattr__(self, name): return self.events.get(name, self.noop)
— event_manager.py:117-121。
这行代码解释了 lfx 为什么能脱离 Web 界面独立运行。create_stream_tokens_event_manager(event_manager.py:139)比默认版少注册了 on_vertices_sorted 和 on_remove_message;组件代码里照样敢直接写 event_manager.on_vertices_sorted(...)——落到 __getattr__ 变成 noop,什么也不发生。组件不需要知道自己跑在哪种运行时里。
下面这段示意代码把这个模式提炼出来(# 示意,非源码):
# 示意,非源码:为什么 noop 兜底能让同一份组件代码跑在不同运 行时里
class Emitter:
def __init__(self, registered):
self.events = registered # 这个运行时支持哪些事件
def _noop(self, *, data): pass
def __getattr__(self, name): # 没登记的事件名 → 静默丢弃
return self.events.get(name, self._noop)
em = Emitter({"on_token": print})
em.on_token(data="hi") # 有登记 → 真的发
em.on_vertices_sorted(data=[]) # 没登记 → 什么都不做,也不报错
重点看:调用方永远不用先判断"这个运行时支不支持这个事件"。
6. 构建结果的三个去向
同一次点构建,结果会分头去三个地方。别把它们搞混:
| 去向 | 函数 | 目的 | 位置 |
|---|---|---|---|
| 给浏览器 | event_manager.on_end_vertex | 实时更新画布 | src/backend/base/langflow/api/build.py:867 |
| 给数据库 | log_vertex_build | 事后查历史构建 | src/lfx/src/lfx/graph/utils.py:348 |
| 给 webhook 订阅者 | emit_vertex_build_event | webhook 触发时也让 UI 有动画 | src/lfx/src/lfx/graph/utils.py:122 |
log_vertex_build 是纯落库,且有两级降级:设置里没开 vertex_builds_storage_enabled 直接返回;开了 telemetry writer 就把行丢进它的磁盘 outbox(graph/utils.py:413-416),避开请求线程的数据库连接池;writer 没跑起来才直接写库。最外层还包了一个 except Exception: logger.warning(graph/utils.py:441)——记录构建历史失败,绝不影响构建本身。
emit_vertex_build_event 走的是另一条完全独立的通道:webhook_event_manager。第一件事就是 if not webhook_event_manager.has_listeners(flow_id_str): return(graph/utils.py:156)——没人订阅就一个字节都不生产。它由 Graph._execute_tasks 在算出 next_runnable_vertices 之后才调用(graph/graph/base.py:2412),因为 payload 里需要 next_vertices_ids;graph/utils.py:425-427 的注释专门解释了为什么不能在 log_vertex_build 里顺手发。
它也有个不显眼的 except ImportError: pass(graph/utils.py:202):lfx 可以在没装 langflow 的环境里独立使用,那时 webhook_event_manager 根本不存在。
7. 前端:从字节流到节点变色
链路的另一半在 src/frontend/src/utils/buildUtils.ts。
7.1 主入口与降级
buildFlowVerticesWithFallback(buildUtils.ts:160-181)只做一件事:先按配置的投递方式试,撞上"流式不可用"就换轮询重来一次。
try {
return await buildFlowVertices({ ...params });
} catch (e) {
if (e.message === POLLING_MESSAGES.ENDPOINT_NOT_AVAILABLE ||
e.message === POLLING_MESSAGES.STREAMING_NOT_SUPPORTED) {
return await buildFlowVertices({ ...params, eventDelivery: EventDeliveryType.POLLING });
}
throw e;
}
两个哨兵消息定义在 src/frontend/src/constants/constants.ts:935-938。这层降级是给那些会缓冲响应体、把长连接一次性吞掉的反向代理准备的。
buildFlowVertices(buildUtils.ts:210)本体按模式分三条路:
direct:直接对 POST/build做流式读(buildUtils.ts:282-328),不拿job_id。streaming:先fetch(buildUrl)拿job_id(buildUtils.ts:336-361),再对事件 URL 做流式读(:385-421)。polling:拿到job_id后交给pollBuildEvents(:431)。
轮询实现在 src/frontend/src/customization/utils/custom-poll-build-events.ts:循环 GET,把 NDJSON 按行 JSON.parse,单行解析失败只跳过这一行(:66-68),看到 end 事件就停,否则 BUILD_POLLING_INTERVAL(25ms,constants.ts:955)后再来。
7.2 解帧:按 \n\n 切,跨 chunk 拼
performStreamingRequest(src/frontend/src/controllers/API/api.tsx:330)负责把字节流切成事件。后端每条事件是 json.dumps(...) + "\n\n"(event_manager.py:87),前端就按 \n\n 切(api.tsx:377)。
TCP chunk 边界和事件边界不对齐,所以有个 current: string[] 缓冲:一段字符串如果不是以 } 结尾、或者拼起来 JSON.parse 还是失败,就先攒着,等下一个 chunk 续上(api.tsx:381-394)。
请求头里那句 Connection: "close" 带着注释——"this flag is fundamental to ensure server stops tasks when client disconnects"(api.tsx:343)。它是第 3.1 节 on_disconnect 那套取消机制在前端这一侧的前提。
7.3 批处理:为什么高频 token 不会卡死 React
这是前端最值得看的一处工程决策。
问题:一次流式对话每秒能产上百条 token 事件,每条都同步 set() 一次 Zustand,React 就要重渲染上百次,界面直接卡住。
做法:把事件 分成两类。
export const BATCHABLE_EVENTS = new Set(["end_vertex", "build_start", "build_end"]);
export const BATCH_YIELD_MS = 50;
— buildUtils.ts:59-65。
processBatchedEvents(buildUtils.ts:606-668)的处理逻辑:
- 属于
BATCHABLE_EVENTS的 → 同步处理,不await。因为一旦await,就跨过了微任务边界,React 18 的自动批处理就断了,同一个 chunk 里的多次set()会变成多次渲染。 - 不属于的(
token/add_message/error/log…)→ 老老实实await onEventFallback(...),它们有真正的异步副作用。 - 这批里只要有过 batchable 事件,最后
await new Promise(r => setTimeout(r, BATCH_YIELD_MS))(buildUtils.ts:665)主动让出 50ms,给浏览器一次真正提交渲染、响应交互的机会。
调度侧的配合在 api.tsx:398-403:一个 chunk 解出来的全部事件先攒成数组,整批交给 onDataBatch,而不是一条条 onData。
// 示意,非源码:批处理与逐条处理的渲染次数差别
// 逐条:await 每条 → 每条一次 set → N 次渲染
for (const e of events) { await handle(e); }
// 成批:同类事件同步处理 → React 合并成 1 次渲染 → 再让出 50ms
for (const e of events) { if (batchable(e)) handleSync(e); else await handle(e); }
await sleep(50);
重点看:同步与否决定了 React 会不会把多次状态更新合并成一次渲染。
7.4 end_vertex:一条事件干了四件事
processEndVertexEvent(buildUtils.ts:456-596)被单独抽成同步函数,注释写明了原因:轮询模式要能不 await 地调它,否则批处理失效。它依次做:
- 判定成败:
buildData.valid为假就从outputs里把错误信息掏出来,调onBuildError,状态置ERROR;否则置BUILT(:474-503)。 - 算分段耗时:如果这个点是输出类节点,用
Date.now() - buildStartTime算出这一段的耗时,写进最后一条机器消息的build_duration,同时更新 React Query 缓存、Zustand、以及持久化到后端(:505-567)。 - 点亮下一批边:
clearAndSetEdgesRunning(nextIds)一次遍历里把所有边的动画清掉,再把源头在nextIds里的边设成animated + className:"running"(:575,实现在src/frontend/src/stores/flowStore.ts:1196-1209)。这就是画布上那 圈流动的虚线。 - 推进层次:把
next_vertices_ids标成TO_BUILD并回调onBuildStart(:581-593)。
7.5 状态怎么落到节点上
useFlowStore(src/frontend/src/stores/flowStore.ts)持有构建状态。关键几处:
| 做什么 | 符号 | 位置 |
|---|---|---|
| 发起构建、装配所有回调 | buildFlow | flowStore.ts:810 |
收 end_vertex 的主回调 | handleBuildUpdate | flowStore.ts:963 |
| 节点 id → 构建状态 + 时间戳 | updateBuildStatus | flowStore.ts:1246 |
| 节点 id → 构建产物(结果/日志) | addDataToFlowPool | flowStore.ts:234 |
| 边的动画开关 | clearAndSetEdgesRunning | flowStore.ts:1196 |
| 收尾把残留的 BUILDING 改回 BUILT | revertBuiltStatusFromBuilding | flowStore.ts:1259 |
handleBuildUpdate 里有个容易忽略的动作:把 next_vertices_ids 和 top_level_vertices 先 zip 再过滤掉已在上一层出现过的,才追加成新的一层(flowStore.ts:989-1015)。因为后端是"每个点各自报告后继",同一个后继会被多个前驱各报一次,前端必须自己去重,否则层列表会膨胀。
7.6 旧路径:updateVerticesOrder
updateVerticesOrder(buildUtils.ts:99-158)先调 getVerticesOrder 拿到排序,再自己一层层 POST 每个点——就是同文件里那个 buildVertices(buildUtils.ts:853)配套的旧模型。它现在被事件流路径取代了:新路径里同样的信息由 vertices_sorted 事件送达(buildUtils.ts:722-748),少一次往返。对应的后端路由 POST /build/{flow_id}/vertices/{vertex_id} 已标 deprecated=True, include_in_schema=False(chat.py:460)。
8. 取消与断线:三条路径通向同一个 cancel
| 触发 | 前端动作 | 后端动作 |
|---|---|---|
| 用户点停止 | buildController.abort() → abort 监听器 POST /build/{job}/cancel | cancel_flow_build → queue_service.cancel_job |
| 关标签页 / 断网 | 连接断开 | on_disconnect → event_task.cancel() |
| 某个组件构建失败 | onBuildError 里 buildController.abort() | 同第一行 |
前端把取消绑在 AbortController 的 abort 事件上(buildUtils.ts:366-379):一次 abort 既停掉本地的读流,又顺手发出取消请求。
cancel_flow_build(build.py:1035-1117)的返回值语义值得留意——它区分了三种"没取消成":
event_task.done()→ 返回True。没得取消也算成功(build.py:1089-1091)。event_task is None且后端支持跨 worker 取消 → 通过 Redis 发信号,返回True(build.py:1064-1078)。event_task is None且不支持 → 返回False,注释诚实承认:"任务已经跑完被清理了"和"任务在一个够不着的 worker 上"这两种情况无法廉价区分(build.py:1079-1087)。
取消路径上还有一处细节:_run_vertex_build 捕到 CancelledError 时,用 asyncio.create_task 而不是 background_tasks.add_task 去收尾 trace(build.py:945-952)。原因写在注释里——POST /build 的响应早就返回了,FastAPI 的 background_tasks 队列已经排干,这时 add_task 会被静默丢弃。
9. 非交互入口:/run 和 webhook 有什么不同
同一台图执行引擎,还有两个不经过 job 队列的入口。
POST /build/{flow_id}/flow | POST /run/{flow_id_or_name} | POST /webhook/{flow_id_or_name} | |
|---|---|---|---|
| 谁在用 | 画布 / Playground | 外部程序、SDK | 外部系统推事件 |
| 认证 | 会话用户 | API key | webhook 认证 |
| 返回 | job_id,结果走事件流 | RunResponse(全部输出)或 SSE | 立刻 202,什么都不返回 |
| 驱动方式 | 递归 spawn(4.2 节) | Graph.arun → process 层批 | 同 /run |
| 事件 | 10 种全套 | 仅 stream=true 时 4 种 | 仅当 UI 开着 SSE 时 |
| 入口 | api/v1/chat.py:256 | api/v1/endpoints.py:962 | api/v1/endpoints.py:1194 |
/run(simplified_run_flow,endpoints.py:962)走 _run_flow_internal(endpoints.py:798)。分两支:
stream=false(默认)→await simple_run_flow(...),一路Graph.from_payload→graph.set_run_id→run_graph_internal→Graph.arun(src/lfx/src/lfx/graph/graph/base.py:1212),同步等到全部跑完,返回RunResponse(outputs=..., session_id=...)(endpoints.py:494)。stream=true→ 自建一个asyncio.Queue和create_stream_tokens_event_manager,把run_flow_generator丢进create_task,返回text/event-stream(endpoints.py:858-881)。注意它用的是精简版事件管理器(event_manager.py:139),只有add_message/token/end/end_vertex/error/build_start/build_end/log,没有vertices_sorted——因为没有画布要涂色。
webhook(webhook_run_flow,endpoints.py:1194)更彻底:把请求体塞成 webhook 组件的 tweak,asyncio.create_task 起一个后台任务,立刻返回 202 {"status": "in progress"}(endpoints.py:1262-1281)。调用方拿不到任何结果。
那 webhook 触发时,开着画布的用户为什么还能看到节点动?靠一条独立的 SSE 通道:GET /webhook-events/{flow_id_or_name}(endpoints.py:1115)订阅 webhook_event_manager,超时就发 heartbeat(endpoints.py:1161-1175)。webhook 任务只在 has_ui_listeners 为真时才 emit_events=True(endpoints.py:1257,1271)——没人看就不生产事件。
10. 边界与坑
event_delivery的默认值前后端不一致。 后端build_flow签名里默认POLLING(chat.py:266),get_build_events默认STREAMING(chat.py:400),前端buildFlow默认STREAMING(flowStore.ts:818)。实际生效的是前端显式带上的那个查询参数。job_id只是内存里的键。 内存版JobQueueService把队列存在进程字典里(services/job_queue/service.py:97)。多 worker 且没配 Redis 后端时,事件 GET 落到别的 worker 上就是 404。跨 worker 能力靠RedisJobQueueService(service.py:719)。- 构建队列是无界的。
create_queue里asyncio.Queue()没有maxsize(service.py:198),所以send_event的QueueFull分支在默认配置下不会触发——它是为带背压的后端准备的。 - 一次未被消费的构建也会跑完。 只要没人断线、没人取消,
generate_flow_events就一直跑,事件在队列里堆着。清理靠JobQueueService的周期清理 + 5 分钟宽限期(service.py:103)。 - 轮询模式下 token 的实时感取决于轮询间隔。 25ms 已经很密,但和流式推送不是一回事。
- 代码里看不出每种事件的正式 schema 定义——事件 payload 是各调用点手写的 dict,没有集中的 pydantic 模型约束。想知道某个事件带什么字段,只能去发送点看。
11. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 建 job、返回 job_id | src/backend/base/langflow/api/build.py | start_flow_build |
| 事件投递分岔(流式/轮询/直连) | src/backend/base/langflow/api/build.py | get_flow_events_response |
| 流式响应 + 心跳 + 断线取消 | src/backend/base/langflow/api/build.py | create_flow_response, _heartbeat, on_disconnect |
| 构建主流程 | src/backend/base/langflow/api/build.py | generate_flow_events |
| 建图与排序 | src/backend/base/langflow/api/build.py | build_graph_and_get_order, create_graph, sort_vertices |
| 单点执行与错误包裹 | src/backend/base/langflow/api/build.py | _build_vertex |
| 递归 spawn 并发 | src/backend/base/langflow/api/build.py | build_vertices, _run_vertex_build |
| 取消构建 | src/backend/base/langflow/api/build.py | cancel_flow_build |
| 断线检测的响应类 | src/backend/base/langflow/api/disconnect.py | DisconnectHandlerStreamingResponse |
| 三个路由 | src/backend/base/langflow/api/v1/chat.py | build_flow, get_build_events, cancel_build |
| 队列服务与崩溃兜底 | src/backend/base/langflow/services/job_queue/service.py | create_queue, start_job, _guarded_task |
| 事件对象与三层容错 | src/lfx/src/lfx/events/event_manager.py | EventManager, send_event, register_event, __getattr__ |
| 事件清单 | src/lfx/src/lfx/events/event_manager.py | create_default_event_manager, create_stream_tokens_event_manager |
| 构建记录落库 | src/lfx/src/lfx/graph/utils.py | log_vertex_build |
| webhook 实时事件 | src/lfx/src/lfx/graph/utils.py | emit_vertex_build_event, emit_build_start_event |
| 单点构建 / 层批驱动 | src/lfx/src/lfx/graph/graph/base.py | build_vertex, process, _execute_tasks, arun |
| token 跨线程发送 | src/lfx/src/lfx/custom/custom_component/component.py | _process_chunk, set_event_manager |
| 前端主入口与降级 | src/frontend/src/utils/buildUtils.ts | buildFlowVerticesWithFallback, buildFlowVertices |
| 事件分发 | src/frontend/src/utils/buildUtils.ts | onEvent, processEndVertexEvent |
| 批处理与让出 | src/frontend/src/utils/buildUtils.ts | processBatchedEvents, BATCHABLE_EVENTS, BATCH_YIELD_MS |
| 流解帧 | src/frontend/src/controllers/API/api.tsx | performStreamingRequest |
| 轮询循环 | src/frontend/src/customization/utils/custom-poll-build-events.ts | customPollBuildEvents |
| 构建状态落到画布 | src/frontend/src/stores/flowStore.ts | buildFlow, handleBuildUpdate, updateBuildStatus, clearAndSetEdgesRunning |
| 非交互入口 | src/backend/base/langflow/api/v1/endpoints.py | simplified_run_flow, simple_run_flow, webhook_run_flow, webhook_events_stream |
继续读: 图是怎么建出来的看 02-graph-construction;get_next_runnable_vertices 凭什么说某个点可跑,看 03-scheduler;后继里出现环怎么办,看 04-cycles-and-branching;一条流怎么变成别人能调的工具,看 06-ecosystem-and-exits。