Langflow 持久化后台执行与人机协同(HITL)暂停-恢复机制深度解析
Langflow 的 Durable HITL(Human-in-the-Loop)机制允许一个 Flow 在运行中途暂停、等待人工审批或输入,再把执行状态完整写库检查点化(checkpoint),即使进程重启也能从断点精确恢复,且不会重复执行已完成的节点。本文基于仓库内的特性规格文档 durable-execution-hitl.md 展开,完整覆盖其领域模型、行为规格、架构决策记录(ADR)与技术实现细节,并结合 api/build.py、jobs/service.py 等源码逐层印证暂停-恢复的底层调用链。读完本文,你将掌握:暂停信号如何建模、单次恢复(single-flight resume)如何通过原子状态翻转实现、恢复为何不等于重跑、以及整个运行如何在后端 Trace 中保持为一条完整链路。
一、功能总览与业务背景
核心命题
HITL 让 Flow 在运行中途暂停,等待人类批准(approve)、拒绝(reject)或提供输入,然后从暂停处精确恢复。关键在于暂停是持久化(durable)的:状态检查点写入数据库,进程重启后依然可恢复,且恢复时不会重新执行、更不会重新计费已完成的 LLM 调用。暂停-恢复的整个运行在后端 Trace 中可完整观测,审批步骤和最终输出都被记录为真实的 span。
业务场景
Agentic Flow 会执行有实际后果的动作(调用工具、发送请求、消耗 token)。团队需要一个人工控制点:在高危步骤执行前由人把关,同时保证浏览器关闭或服务重启不会丢失整个运行。HITL 将 Flow 从“一次性的 fire-and-forget 请求”变成了长生命周期、可恢复的工作流。
限界上下文
- 所属上下文:Flow Execution — 持久化后台任务与图检查点(Durable Background Jobs & Graph Checkpointing)。
- 关联上下文:
- Graph Engine(客户-供应商):在图层边界处抛出并恢复暂停;
- Background Execution / Job Service(合作):持久化任务状态,提供单次恢复(single-flight resume);
- Tracing(跟随者):将整个运行——包括审批闸门与最终输出——记录为 span;
- 前端 AG-UI(客户-供应商):重新挂接到实时事件流并渲染决策界面。
默认关闭(Default-off)保证
暂停探测(pause probe)在没有任何组件请求暂停时是 no-op:普通 Flow 完全不受影响——不写检查点、不产生额外 span、不引起状态变化。这意味着该特性可以安全上线。
二、通用语言(Ubiquitous Language)
文档为整个 HITL 领域定义了统一术语表,每一项都锚定了具体的代码符号,读者可以按下表快速定位实现:
| 术语 | 定义 | 代码引用 |
|---|---|---|
| Pause request | 组件发出的“继续前需要人工决策”信号 | PauseRequested、GraphPausedException |
| Checkpoint | 图层边界处序列化的图状态快照,写入数据库 | lfx/graph/checkpoint/schema.py、checkpoint 表 |
| Suspend | 任务转入持久化等待状态 | JobStatus.SUSPENDED |
| Resume | 从检查点重新水合(re-hydrate)图并越过暂停继续 | resume_from_checkpoint、build_resumed_graph_and_get_order |
| Single-flight resume | 原子的 SUSPENDED → IN_PROGRESS 认领,保证决策只应用一次 |
claim_suspended_for_resume |
| Human input request | 展示给用户的待答请求(工具审批或 HumanInput 节点) | get_pending_human_request、human_input_required 事件 |
| Decision | 人类的回答:action_id(如 approve/reject)加可选 values |
graph.human_input_decisions |
| Run id | 图的追踪身份;等于消息上的 graph_run_id |
graph.set_run_id、trace id |
| Job id | 用于检查点与恢复的持久化任务身份 | JobStatus、/api/v2/workflows/{job_id}/resume |
| Gate span | 记录已解析决策的追踪 span:“Human In The Loop — {动作标签}”(如 Approve/Reject/Remove) | TracingService.record_event_span |
三、领域模型
3.1 Durable Run 聚合体
- 根实体:后台 Job(
job_id); - 实体:Graph(携带
run_id)、Checkpoint、Pending Human Request; - 值对象:
JobStatus、Decision(action_id+values)、request_id(格式为{node_id}:{run_id};Agent 工具审批会追加 LangGraph 的 interrupt id,即{node_id}:{run_id}:{interrupt_id}——这使得一次运行中的每次审批都可单独寻址,且针对第 N 次审批的过期 resume 会在第 N+1 次审批期间被拒绝)。
这一 request_id 规则在源码中可以直接印证,见 hitl.py 的 request_id_targets_vertex:
def request_id_targets_vertex(request_id: str, vertex_id: str, run_id: str) -> bool:
"""Whether a resume ``request_id`` addresses this vertex's pause.
Node pauses use ``vertex_id:run_id``; agent tool-approval pauses append a
per-pause interrupt id (``vertex_id:run_id:interrupt_id``) so each approval
in one run is individually addressable and a stale resume for an earlier
approval cannot be applied to a later one.
"""
base = f"{vertex_id}:{run_id}"
return request_id == base or request_id.startswith(base + ":")
聚合不变式(Invariants):
- 暂停必须先持久化检查点,运行才能挂起(不允许丢失状态);
- 恢复必须是单次(single-flight)的:一个决策恰好应用一次;
- 恢复不得重新执行已构建(built)的顶点(不重复计费、不重复触发工具);已死的分支保持死亡;
- 整个运行(暂停前 + 恢复后)必须落在同一个
run_id的追踪中。
3.2 领域事件
| 事件 | 触发时机 | 负载 | 消费者 |
|---|---|---|---|
human_input_required |
组件抛出暂停 | 卡片/请求描述 | 前端(渲染卡片/栏);任务运行器(挂起) |
Job SUSPENDED |
检查点已写入、运行挂起 | job_id、request_id |
待处理请求 API;UI 徽章 |
| Resume accepted | 有效决策被认领 | job_id、decision |
图重建;追踪 |
| Gate resolved span | 恢复时重新初始化追踪 | “Human In The Loop — {decision}” | 追踪面板(/monitor/traces/{id}) |
Job COMPLETED / 终态 span |
恢复后的运行结束 | Chat Output span + flush | 追踪面板;消息历史 |
四、行为规格:暂停、恢复与边界场景
核心特性以用户故事表述:作为 Flow 作者,我希望某个步骤暂停等待人工审批并能在中断后存活,以便高危动作被把关、服务重启不丢运行。前置条件:Flow 包含可请求审批的 Agent 工具(或 HumanInput 节点),且启用了持久化后台执行。
场景 1:Agent 工具暂停等待审批
Agent 决定调用一个受控工具 → 图写检查点 → 任务变为 SUSPENDED → human_input_required 事件弹出 Approve/Reject 卡片。
场景 2:恢复可穿越进程重启
任务处于 SUSPENDED 且检查点已落盘;后端进程重启、用户点击批准后,图从检查点重建并越过暂停点继续。
场景 3:恢复不是重跑
- 已构建的顶点被恢复而非重新执行(无 LLM 重复计费、无工具重复触发);
- 恢复-已构建的顶点永远不会被 runnable-predecessor 回溯交回构建循环(否则重跑的 Chat Input 会为同一轮持久化一条重复的
User消息); - 暂停前已被判死的分支保持死亡。
场景 4:嵌套 Flow 中的 HITL 被拒绝
若 Run Flow / Sub Flow / flow-as-tool 的目标 Flow 内含已连线的 Human Input 节点(或受审批闸门的 Agent 工具),运行会显式失败并给出清晰错误(“cannot run as a nested flow… 请把审批移到父 Flow”),而不是悄悄不暂停。但若目标 Flow 中的 Human Input 未接线(无下游消费者),运行不受影响——它在运行时被跳过而非阻塞。
场景 5:两个 HITL 节点串行
第一个 HumanInput 把关第二个的路径。第一个批准后,第二个暂停;第二个再批准后,运行只终结一次(不会绕回第一个节点),且每个节点只有被选中的分支会执行——第一个节点未被选中的分支在第二次暂停的恢复中依然保持死亡。
场景 6:重跑 Flow 会取代(supersede)过期暂停
- 同一用户对同一 Flow 提交新的后台运行时,过期的
SUSPENDED运行先被取消(CANCELLED、清理 pending 请求),新运行才开始; - 被取代暂停的持久化聊天卡片被打上
superseded标记,历史重载后渲染为关闭态(“Superseded by a new run”),而不会重新变为可交互、每次点击都 409; - 作用域:langflow 后端按 flow + user(另一用户在同一 Flow 上的暂停不受影响);
lfx serve按 flow(单一 API-key 身份);运行中的任务永不被取代,并行运行保持支持。
场景 7:Supersede 输给并发的 Resume
若恢复请求抢先赢得原子的 SUSPENDED → IN_PROGRESS 认领(发生在 supersede 的 select 与 cancel 之间),supersede 完全跳过该任务——不写 CANCELLED、不删检查点、不清元数据(后端)或不发残留 STOP 信号(lfx serve),被恢复的运行正常完成。
场景 8:单次恢复(Single-flight resume)
同一 SUSPENDED 任务收到两个恢复请求时,恰好一个赢得 SUSPENDED → IN_PROGRESS 翻转;另一个被拒绝(NOT_RESUMABLE)。
场景 9:恢复后的运行是一条完整的后端 Trace
单个 trace 完整持有 Chat Input → Human In The Loop — {decision} → … → Chat Output;追踪面板在整页刷新后仍展示闸门与终态输出,且没有任何客户端侧持久化。
场景 10:超时后的决策被改道而非丢失
对配置了超时策略的待处理请求,若决策在截止后到达,reroute_decision_on_timeout 会确定性地解析出有效决策。源码实现见 hitl.py:采用惰性超时(无后台看门狗)——人类响应时比较当前时间与 paused_at + timeout_seconds;已过期则改道到节点配置的 fallback 分支,否则落入 __expired__ 哨兵动作(不选任何分支,HumanInput.route_branch 使所有分支停止)。
五、架构决策记录(ADR)
ADR-001:暂停是信号,不是失败(已接受)
背景:需要人工输入的组件必须在不破坏状态、不把运行标记为失败的前提下停下图。
决策:将暂停建模为专用异常 GraphPausedException,在图层边界抛出——此时图已完成自我快照并写入检查点。构建缝隙(build seam)捕获它并发出非终态的 human_input_required 事件——事件流在没有 on_end 的情况下结束,因此运行是“等待中”而非“已完成”。源码注释明确强调这是控制流而非错误:“callers must let it propagate unwrapped past generic error handlers so the run is finalized as suspended, never failed.”
后果:
- 收益:“等待”与“失败”清晰分离;天生可恢复;默认关闭(无探测 → no-op);
- 代价:每个图层边界都要咨询一次暂停探测(开销极低;无待处理暂停时完全跳过)。
ADR-002:恢复是单次的且可回滚(已接受)
背景:即使存在重复点击、重试或并发调用者,决策也必须恰好应用一次。
决策:用原子的 SUSPENDED → IN_PROGRESS 状态翻转(claim_suspended_for_resume)认领任务;恢复过程中若失败则回滚状态,决策永远不会被静默吞掉。
实现证据:jobs/service.py 的 claim_suspended_for_resume 正是用一条条件 UPDATE ... WHERE status = SUSPENDED 完成——“只有 rowcount == 1 的并发恢复者获胜”。直达 IN_PROGRESS(并附带新心跳)还能让并发的启动清扫看到非 QUEUED 的活动行,避免重复入队。
后果:
- 收益:决策 exactly-once;重试安全;审批不丢失;
- 代价:过期/重复的恢复返回
NOT_RESUMABLE(由 UI 处理)。
ADR-003:恢复 ≠ 重跑(已接受)
背景:恢复时重新执行已完成的顶点会重复计费 LLM、重复触发工具,甚至可能复活已死的分支。
决策:从检查点恢复已构建顶点;只重跑被暂停顶点的非输入前置节点中、其丢失输出确实会被未构建消费者读取的那些(_rerun_non_input_predecessors / _unbuild_needed_dropped_producers)。位于“仍然已构建、已往返”的消费者之后的生产者保持不动——重跑它是浪费,且对带副作用的节点(Agent 会重新计费 LLM 并重复发消息)会在之后的每次恢复中表现为重复输出。输入类节点(如 Chat Input)从不重跑。
关键守卫:仅选定恢复层不足以守住该不变式。恢复会原样从检查点恢复 vertices_to_run,因此以已构建状态回来的顶点仍留在可运行池里;由于 RunnableVerticesManager.is_vertex_runnable 从不咨询 built 标志,当后继被阻塞时寻找可运行前置的后向遍历(find_runnable_predecessors_for_successor)可能在恢复中途把已构建顶点交回构建循环。因此 Graph.is_vertex_runnable 会拒绝那些仍为 built 且从检查点恢复(checkpoint_restored_built_ids)的顶点——但排除循环顶点(它们本就应当重跑)。读取实时的 built 标志而非检查点中的标志,正是这一点让“被闸门的节点”和“不透明的被丢弃生产者”(在检查点时已构建、恢复时故意解构建)保持可选。
后果:
- 收益:无双重计费、无重复副作用、确定性的继续执行;
- 代价:恢复时重跑一个最小的前置集合;终态输出在恢复时新鲜产出(见 ADR-004);
- 预防的症状:被重执行的 Chat Input 会为同一轮持久化第二条
User消息,使暂停中的聊天在决策后把用户问题渲染两次。
相关单元测试(如 test_resume_does_not_reschedule_built.py)专门守护这一不变式。
ADR-004:整个运行是一条持久化的后端 Trace(闸门 + Chat Output 作为 span)(已接受,2026-06-23)
背景:早期实现中,恢复后的运行在后端丢失追踪数据——Chat Output span 缺失,Human In The Loop 步骤只存在于前端注入、持久化于 localStorage 的节点中:其他设备/用户看不到,数据库里也不完整。根因有二:恢复路径从未初始化追踪(trace_context_var 未设置,导致暂停后的顶点跳过 add_trace);初始运行用一次性 uuid 而非 run_id 追踪,把运行劈成了孤儿 trace。
决策:
- 恢复时调用
graph.initialize_run(),传入检查点中的run_id,使恢复后的顶点追踪进同一条 trace,并恢复flow_name用于 trace 标题; - 在
build_graph_from_data中于initialize_run之前钉住调用方的run_id(由create_graph转发),使初始运行追踪进graph_run_id而非新 uuid; - 通过
TracingService.record_event_span把已解析的闸门记录为真实 span(“Human In The Loop — {动作标签}”,如 Approve/Reject/Remove); - 移除前端
localStorage持久化;TraceDetailView从后端 trace 读取闸门与输出,只在实际时间窗内合成闸门节点,并按名称去重。
后果:
- 收益:一条完整 trace(
Chat Input → 闸门 → … → Chat Output)在刷新、跨设备、跨用户间持久;没有客户端侧补丁; - 代价:恢复时重新初始化一次追踪(追踪关闭时是 no-op);
- 产品影响:追踪面板成为 HITL 运行行为的唯一事实来源。
实现证据:api/build.py 的 build_resumed_graph_and_get_order 完整呈现了这条决策链:
async def build_resumed_graph_and_get_order() -> tuple[list[str], list[str], Graph]:
"""Resume a suspended HITL run from its durable checkpoint instead of building fresh.
Hydrates the graph, injects the human decision keyed by request_id, and un-builds
the paused node so it re-runs and routes (its first-run output was a placeholder).
"""
...
graph = LfxGraph.resume_from_checkpoint(checkpoint, checkpoint_store=store)
if not graph.user_id:
graph.user_id = str(current_user.id)
# Resume skips the initial run's trace setup (trace_context_var stays unset → post-pause
# vertices like Chat Output never trace); re-init so the resumed vertices trace.
graph.flow_name = graph.flow_name or flow_name
await graph.initialize_run()
pending = await get_job_service().get_pending_human_request(job_id)
decision = reroute_decision_on_timeout(pending, resume["decision"])
# Merge with checkpoint-restored decisions so a re-run HITL keeps its answer (no multi-HITL loop).
graph.human_input_decisions = {
**(getattr(graph, "human_input_decisions", {}) or {}),
resume["request_id"]: decision,
}
action_id = str((decision or {}).get("action_id", ""))
gate_label = _hitl_gate_label(action_id, (pending or {}).get("options"))
if graph.tracing_service:
graph.tracing_service.record_event_span(
span_id=f"hitl-{resume['request_id']}",
name=f"Human In The Loop — {gate_label}",
outputs={"decision": action_id},
)
for vertex in graph.vertices:
if request_id_targets_vertex(str(resume["request_id"]), vertex.id, run_id):
vertex.built = False
_rerun_non_input_predecessors(graph, vertex.id)
first_layer = graph.resume_first_layer()
...
return first_layer, list(graph.vertices_to_run), graph
注意两处细节:决策以合并方式写入 graph.human_input_decisions(保留检查点恢复的历史决策,避免多 HITL 循环);HITL 动作是用户自定义的(Approve/Reject/Remove/…),所以 gate span 的标签取自待处理请求选项中的按钮标签,而非硬编码的 approve/reject 二元组。
ADR-005:终态 span 冲刷竞态修复(已接受)
背景:组件 span 通过 worker 队列异步记录。运行结束时,worker 在终态组件的 end 事件仍在队列中时就被取消,导致约 8 次运行中 7 次丢失 Chat Output span。
决策:取消 worker 后内联排空队列(service.py::_stop),并在 flush 时强制完成“已开始但未结束”的 span(native.py::_finalize_pending_spans)。
后果:
- 收益:任何已执行的组件都不会被静默地从 trace 中丢弃;
- 代价:关闭时一次微小的同步排空(受队列大小约束)。
六、技术规范:端到端调用链
6.1 提交时取代过期暂停(Submit supersedes stale pauses)
BackgroundExecutionService.submit在create_job之后、入队之前调用supersede_suspended_runs(flow_id, user_id)(放在create_job之后,保证幂等重试仍能返回既有任务):同一 flow + user 的每个SUSPENDED工作流任务都经_cancel_suspended路径取消(状态CANCELLED、卡片标记 superseded、检查点删除、pending 清理、run_cancelled事件、总线关闭)。- 取消是一个原子认领:
JobService.claim_suspended_for_cancel(条件UPDATE … WHERE status = SUSPENDED,与claim_suspended_for_resume互为镜像)翻转该行,清理只针对本调用方实际认领的行。抢先赢得翻转的 resume 保留其运行——supersede 永远不能取消已恢复的任务或在中途销毁其检查点。stop_job使用同样的认领,失败时落回运行中任务的 STOP 路径。源码见 claim_suspended_for_cancel。 - 清除元数据之前,
mark_card_superseded(api/v2/hitl.py)把持久化卡片的human_input内容补丁为superseded: true(跳过已回答的卡片);HumanInputCard将该状态渲染为关闭态,因此重载的历史无法再提供一个只会 409 的决策。 DurableServeWorkflowHost.submit_background按 flow 作用域镜像同一逻辑:先SqliteDurableJobStore.claim_suspended_for_cancel,再(仅对认领到的行)发 STOP 信号、取消任务、清 pending 元数据——认领失败同样意味着不给恢复中的继续执行留下游离 STOP 信号。- 前端:画布徽章自动打开以
request_id为键,因此取代运行的新暂停会重新打开曾被关掉的弹出层;5 秒一次的 pending 轮询使所有界面收敛到唯一剩下的暂停。
后端行为有对应测试覆盖:test_supersede_suspended.py。
6.2 暂停 → 挂起(初始运行)
- 图在每个图层边界咨询暂停探测;发现待处理暂停时自我快照,写入
checkpoint行(始终可写的 schema),然后抛出GraphPausedException; - 构建缝隙(api/build.py)捕获它,把人工输入卡片持久化到历史,发出
human_input_required,并以无on_end结束事件流; - 运行器把任务标记为
SUSPENDED(services/background_execution/runner.py)。
6.3 恢复(单次)
POST /api/v2/workflows/{job_id}/resume,负载为{request_id, decision}(路由定义见 workflow_execution.py);- 路由在应用决策前重新强制
FlowAction.EXECUTE权限(ensure_resume_execute_permission)——共享访问被撤销后,仅凭任务归属权不够; claim_suspended_for_resume执行原子状态认领;失败 →NOT_RESUMABLE;- 已回答卡片通过
mark_card_answered打标,使用在入队继续执行之前快照的card_message_id,且仅当卡片的request_id匹配时——第一次决策永远无法解决第二个暂停的卡片; - 被闸门的顶点通过
request_id_targets_vertex(lfx.run.hitl)定位,同时接受节点形态与带 nonce 后缀的工具审批形态; build_resumed_graph_and_get_order(api/build.py):
graph = LfxGraph.resume_from_checkpoint(checkpoint, checkpoint_store=store)
graph.flow_name = graph.flow_name or flow_name
await graph.initialize_run() # 在原 run_id 上重新初始化追踪
decision = reroute_decision_on_timeout(pending, resume["decision"])
graph.human_input_decisions = {resume["request_id"]: decision}
action_id = str((decision or {}).get("action_id", ""))
# HITL 动作是用户自定义的(Approve/Reject/Remove/...),gate span 的标签
# 取自待处理选项中的按钮标签,而非硬编码 approve/reject 二元组。
gate_label = _hitl_gate_label(action_id, (pending or {}).get("options"))
graph.tracing_service.record_event_span(
span_id=f"hitl-{resume['request_id']}",
name=f"Human In The Loop — {gate_label}",
outputs={"decision": action_id},
)
# 恢复已构建顶点;只重跑被闸门顶点的非输入前置
6.4 Trace 连续性(初始运行)
build_graph_from_data(api/utils/flow_utils.py)在追踪启动前钉住运行 id:
if (caller_run_id := kwargs.get("run_id")) is not None:
graph.set_run_id(caller_run_id)
await graph.initialize_run()
create_graph 把 run_id=str(job_id) if job_id is not None else run_id 同时转发给 build_graph_from_data 与 build_graph_from_db。从源码结构看,这一步是整个“一条 trace”不变式的起点:若 run_id 未被钉住,初始运行会追踪进新 uuid,成为孤儿 trace。
6.5 Gate span 记录
TracingService.record_event_span(span_id, name, outputs, span_type="chain") 为每个就绪的 tracer 向活动 trace 写入一对独立的 start+end span;被禁用或无上下文(contextless)的调用是 no-op。
6.6 前端(只读后端 span)
useGetTraceQuery→GET /api/v1/monitor/traces/{id}返回 span 树;TraceDetailView.tsx:渲染后端 gate span;backendHasGate去重防止重复的合成节点;executedOutputSpans只在恢复时间窗内从实时flowPool桥接 Chat Output(按名称去重);hitlStore.ts:纯内存(pending槽位);localStorage/persist已移除;- 已回答卡片的传播:聊天会话缓存是订阅型查询(
staleTime: Infinity,由setQueryData供数),invalidateQueries无法刷新它,只在为空时从后端重新水合。因此在其他界面(画布徽章、trace 栏、另一个标签页)做出的决策,由withAnsweredHumanInputCards(human-input-card.ts)在加载时采纳——把后端副本的submitted_action打到缓存仍显示为打开态的卡片上。否则已回答的暂停会一直渲染按钮,直到整页刷新。
6.7 lfx serve 持久化模式
裸 lfx serve 在没有数据库的情况下获得同样的后台 + HITL 契约,底层是单节点 SQLite 基座:
- 基座(
lfx/services/durable/):SqliteDurableJobStore(任务行、无缺口的逐任务事件日志、控制信号、原子claim_suspended_for_resume)与SqliteCheckpointStore(同一 DB 文件上的既有CheckpointStoreABC)。标准库 SQLite、WAL 模式、经asyncio.to_thread异步化;每次部署一个崩溃安全文件; - Host(serve_durable.py):
DurableServeWorkflowHost扩展ServeWorkflowHost并声明supports_background = True;submit_background创建任务行并派生运行器;get_job_status从存储的output_events重建已完成运行的WorkflowExecutionResponse;stop_job记录持久化 STOP 信号并取消本地任务; - 运行器(
_drive):带检查点地运行graph.process();GraphPausedException→ 任务SUSPENDED+ 待处理请求持久化到任务元数据;完成 → 终态ComponentOutput持久化为任务结果;失败 →FAILED并记录错误; - 恢复:
POST /api/v2/workflows/{job_id}/resume(与后端相同的请求/响应契约,含 404/409 语义以及针对allowed_decisions之外动作的 422INVALID_DECISION守卫)执行单次认领,从检查点恢复——无论进程内还是重启之后——重新应用调用方身份、注入决策、解构建被闸门顶点及其不透明丢弃的生产者,然后重新驱动。Agent 暂停 blob 经set_default_checkpoint_store持久化在同一 SQLite 文件中,因此工具审批中断同样能穿越重启。GET /api/v2/workflows/{job_id}/pending暴露待处理请求。serve 侧的认领逻辑见 serve_durable.py,认领失败时返回“Job is not suspended, already resumed, or the request_id is stale.” - 事件(SSE):
GET /api/v2/workflows/{job_id}/events重新挂接到后台运行——每个任务响应都会广告links.eventsURL。它在Last-Event-ID之后回放持久化事件日志(每帧以seq为键),然后尾随到运行结束或挂起;HITL 暂停会回放并关闭流,客户端在恢复后带Last-Event-ID重连。未知任务 → 404。 - 显式开启(Opt-in):
LFX_SERVE_DURABLE_DB=<path>。未设置时 serve 保持无状态,mode: background保持返回 422;多 worker 模式下的 worker 继承该环境变量并共享 DB 文件(stop 只取消执行该运行的那个 worker 中的运行)。 - 验证证据:test_serve_durable.py ——后台 echo 以同步形态的响应完成;HITL Flow 挂起、暴露其待处理请求、只恢复被批准的分支,并在同一 DB 文件上的全新应用实例上恢复到完成(重启等价性)。
{job_id}/eventsSSE 流回放已完成的运行,在挂起时呈现human_input_request,并在恢复后带Last-Event-ID重连时只投递恢复后的human_input_decision+end帧。
6.8 组件索引(Concern → 代码位置)
| 关注点 | 位置 | 关键符号 |
|---|---|---|
| 暂停信号 + 检查点 | lfx/graph/graph/base.py、lfx/graph/checkpoint/ |
GraphPausedException、resume_from_checkpoint、set_run_id |
| 构建缝隙 + 恢复重建 | api/build.py |
build_resumed_graph_and_get_order、_rerun_non_input_predecessors |
| 恢复调度守卫 | lfx/graph/graph/base.py |
is_vertex_runnable、checkpoint_restored_built_ids |
| Run-id 钉住 | api/utils/flow_utils.py |
build_graph_from_data |
| 持久化任务 + 单次恢复 | services/background_execution/{runner,service}.py |
JobStatus.SUSPENDED、claim_suspended_for_resume |
| lfx serve 持久化基座 | lfx/services/durable/、lfx/cli/serve_durable.py |
SqliteDurableJobStore、SqliteCheckpointStore、DurableServeWorkflowHost |
| 追踪 | services/tracing/{native,service}.py |
record_event_span、_finalize_pending_spans、_stop 排空 |
| Trace UI | TraceComponent/TraceDetailView.tsx |
backendHasGate、executedOutputSpans |
| 实时 HITL 槽位 | stores/hitlStore.ts |
pending(内存) |
七、可观测性
- 追踪面板(Flow Activity):运行渲染为 span 树——
Chat Input → tools → Agent → Chat Output——受闸门运行时额外包含Human In The Loop — {decision}span。经GET /api/v1/monitor/traces/{id}读取; - 任务状态:
SUSPENDED/IN_PROGRESS/COMPLETED/FAILED经 workflows API 暴露;待处理请求经GET /api/v2/workflows/pending?flow_id=查询; - 端到端验证(2026-06-23):一次驱动式的 暂停→批准→恢复 产生一条含 21 个 span 的 trace(
Chat Input → Human In The Loop — Approved → Calculator → URL → Agent → fetch_content → … → Chat Output),已在数据库、/monitor/traces/{id}负载以及整页刷新后的追踪面板中确认。
八、部署与回滚
- 默认关闭:没有组件请求暂停 → 不写检查点、无额外 span、无状态翻动。可安全上线;
- Schema:
checkpoint表与span/trace表支撑持久化与可观测性;本特性不包含破坏性迁移; - 运行模式:在持久化后台路径(
POST /api/v2/workflows,mode: background)下工作,经POST /api/v2/workflows/{job_id}/resume恢复。在裸lfx serve下,同一契约通过LFX_SERVE_DURABLE_DB=<path>显式开启(SQLite 基座,无数据库服务);未设置时 serve 保持无状态; - 回滚:回退 trace 连续性变更会恢复旧行为(Chat Output/闸门缺席于恢复后的后端 trace),但不影响暂停/恢复正确性;回退暂停/恢复层则完全禁用 HITL(Flow 直贯运行)。
九、架构图
9.1 系统上下文(C4 Level 1)
C4Context
title System Context — Durable HITL
Person(user, "User", "Approves/rejects a gated step")
System(langflow, "Langflow", "Flow execution platform")
SystemDb(db, "Database", "checkpoint, job status, trace/span")
Rel(user, langflow, "Runs flow; answers Approve/Reject")
Rel(langflow, db, "Persists checkpoint + status + spans")
Rel(langflow, user, "Surfaces card / trace gate")
9.2 时序:暂停 → 挂起 → 恢复
sequenceDiagram
participant U as User
participant API as Workflows API
participant G as Graph Engine
participant DB as Database
participant T as Tracing
U->>API: POST /v2/workflows (background)
API->>G: build + run (run_id pinned)
G->>T: trace Chat Input … Agent
G->>DB: write checkpoint (layer boundary)
G-->>API: GraphPausedException → human_input_required
API->>DB: job SUSPENDED
API-->>U: Approve / Reject card
U->>API: POST /v2/workflows/{job}/resume {decision}
API->>DB: claim SUSPENDED→IN_PROGRESS (single-flight)
API->>G: resume_from_checkpoint + initialize_run(run_id)
G->>T: record "Human In The Loop — {decision}" span
G->>T: trace re-run predecessors … Chat Output
G->>DB: flush spans (one trace) + job COMPLETED
API-->>U: result; trace panel shows full run
9.3 Trace 连续性(为何是一条 trace)
graph TB
A[POST /v2/workflows] --> B[create_graph]
B --> C[build_graph_from_data]
C --> D{run_id pinned?}
D -->|yes set_run_id job_id| E[initialize_run → trace = graph_run_id]
D -->|no - old bug| F[fresh uuid → orphan trace]
E --> G[pause → checkpoint stores run_id]
G --> H[resume_from_checkpoint set_run_id run_id]
H --> I[initialize_run → SAME trace]
I --> J[record gate span + Chat Output]
J --> K[(One complete trace)]
9.4 前端读取路径(无 localStorage)
graph LR
P[Trace panel] --> Q[useGetTraceQuery]
Q --> R[GET /api/v1/monitor/traces/id]
R --> S[span tree incl. gate + Chat Output]
S --> T[TraceDetailView renders backend spans]
T --> U{backendHasGate?}
U -->|yes| V[use backend gate - no synth]
U -->|live window| W[synth gate / bridge Chat Output from flowPool]
十、延伸阅读与验证入口
想进一步在仓库中追踪本特性的实现与测试,建议按以下入口逐层深入:
- 暂停异常定义:graph/exceptions.py
- 检查点恢复核心(
resume_from_checkpoint、_unbuild_needed_dropped_producers):checkpoint/resume.py - 后端构建缝隙与恢复重建:api/build.py
- 恢复相关单元测试:checkpoint 测试目录,含 test_graph_pause.py、test_resume.py、test_resume_does_not_reschedule_built.py
- 持久化 SQLite 基座与 serve 侧恢复:durable/sqlite_store.py、重启后恢复测试
- 面向用户的 HITL 使用文档:docs/Flows/human-in-the-loop.mdx
适用前提与限制:HITL 暂停-恢复仅适用于持久化后台执行路径(mode: background)或显式开启 LFX_SERVE_DURABLE_DB 的 lfx serve;嵌套 Flow 中已连线的 HITL 节点会被显式拒绝而非静默跳过;决策动作由用户自定义(不限于 approve/reject),超时后的迟到决策会被确定性地改道到 fallback 分支或哨兵动作,而不会被丢弃。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0625
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00