首页
/ Langflow 持久化后台执行与人机协同(HITL)暂停-恢复机制深度解析

Langflow 持久化后台执行与人机协同(HITL)暂停-恢复机制深度解析

2026-09-06 13:33:40作者:冯爽妲Honey

Langflow 的 Durable HITL(Human-in-the-Loop)机制允许一个 Flow 在运行中途暂停、等待人工审批或输入,再把执行状态完整写库检查点化(checkpoint),即使进程重启也能从断点精确恢复,且不会重复执行已完成的节点。本文基于仓库内的特性规格文档 durable-execution-hitl.md 展开,完整覆盖其领域模型、行为规格、架构决策记录(ADR)与技术实现细节,并结合 api/build.pyjobs/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 组件发出的“继续前需要人工决策”信号 PauseRequestedGraphPausedException
Checkpoint 图层边界处序列化的图状态快照,写入数据库 lfx/graph/checkpoint/schema.pycheckpoint
Suspend 任务转入持久化等待状态 JobStatus.SUSPENDED
Resume 从检查点重新水合(re-hydrate)图并越过暂停继续 resume_from_checkpointbuild_resumed_graph_and_get_order
Single-flight resume 原子的 SUSPENDED → IN_PROGRESS 认领,保证决策只应用一次 claim_suspended_for_resume
Human input request 展示给用户的待答请求(工具审批或 HumanInput 节点) get_pending_human_requesthuman_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 聚合体

  • 根实体:后台 Jobjob_id);
  • 实体Graph(携带 run_id)、CheckpointPending Human Request
  • 值对象JobStatusDecisionaction_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)

  1. 暂停必须先持久化检查点,运行才能挂起(不允许丢失状态);
  2. 恢复必须是单次(single-flight)的:一个决策恰好应用一次
  3. 恢复不得重新执行已构建(built)的顶点(不重复计费、不重复触发工具);已死的分支保持死亡;
  4. 整个运行(暂停前 + 恢复后)必须落在同一个 run_id 的追踪中。

3.2 领域事件

事件 触发时机 负载 消费者
human_input_required 组件抛出暂停 卡片/请求描述 前端(渲染卡片/栏);任务运行器(挂起)
Job SUSPENDED 检查点已写入、运行挂起 job_idrequest_id 待处理请求 API;UI 徽章
Resume accepted 有效决策被认领 job_iddecision 图重建;追踪
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 决定调用一个受控工具 → 图写检查点 → 任务变为 SUSPENDEDhuman_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 serveflow(单一 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。

决策

  1. 恢复时调用 graph.initialize_run(),传入检查点中的 run_id,使恢复后的顶点追踪进同一条 trace,并恢复 flow_name 用于 trace 标题;
  2. build_graph_from_data 中于 initialize_run 之前钉住调用方的 run_id(由 create_graph 转发),使初始运行追踪进 graph_run_id 而非新 uuid;
  3. 通过 TracingService.record_event_span 把已解析的闸门记录为真实 span(“Human In The Loop — {动作标签}”,如 Approve/Reject/Remove);
  4. 移除前端 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.submitcreate_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_supersededapi/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 暂停 → 挂起(初始运行)

  1. 图在每个图层边界咨询暂停探测;发现待处理暂停时自我快照,写入 checkpoint 行(始终可写的 schema),然后抛出 GraphPausedException
  2. 构建缝隙(api/build.py)捕获它,把人工输入卡片持久化到历史,发出 human_input_required,并以无 on_end 结束事件流;
  3. 运行器把任务标记为 SUSPENDEDservices/background_execution/runner.py)。

6.3 恢复(单次)

  1. POST /api/v2/workflows/{job_id}/resume,负载为 {request_id, decision}(路由定义见 workflow_execution.py);
  2. 路由在应用决策前重新强制 FlowAction.EXECUTE 权限(ensure_resume_execute_permission)——共享访问被撤销后,仅凭任务归属权不够;
  3. claim_suspended_for_resume 执行原子状态认领;失败 → NOT_RESUMABLE
  4. 已回答卡片通过 mark_card_answered 打标,使用在入队继续执行之前快照的 card_message_id,且仅当卡片的 request_id 匹配时——第一次决策永远无法解决第二个暂停的卡片;
  5. 被闸门的顶点通过 request_id_targets_vertexlfx.run.hitl)定位,同时接受节点形态与带 nonce 后缀的工具审批形态;
  6. build_resumed_graph_and_get_orderapi/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_dataapi/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_graphrun_id=str(job_id) if job_id is not None else run_id 同时转发给 build_graph_from_databuild_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)

  • useGetTraceQueryGET /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 栏、另一个标签页)做出的决策,由 withAnsweredHumanInputCardshuman-input-card.ts)在加载时采纳——把后端副本的 submitted_action 打到缓存仍显示为打开态的卡片上。否则已回答的暂停会一直渲染按钮,直到整页刷新。

6.7 lfx serve 持久化模式

lfx serve 在没有数据库的情况下获得同样的后台 + HITL 契约,底层是单节点 SQLite 基座:

  • 基座lfx/services/durable/):SqliteDurableJobStore(任务行、无缺口的逐任务事件日志、控制信号、原子 claim_suspended_for_resume)与 SqliteCheckpointStore(同一 DB 文件上的既有 CheckpointStore ABC)。标准库 SQLite、WAL 模式、经 asyncio.to_thread 异步化;每次部署一个崩溃安全文件;
  • Hostserve_durable.py):DurableServeWorkflowHost 扩展 ServeWorkflowHost 并声明 supports_background = Truesubmit_background 创建任务行并派生运行器;get_job_status 从存储的 output_events 重建已完成运行的 WorkflowExecutionResponsestop_job 记录持久化 STOP 信号并取消本地任务;
  • 运行器_drive):带检查点地运行 graph.process()GraphPausedException → 任务 SUSPENDED + 待处理请求持久化到任务元数据;完成 → 终态 ComponentOutput 持久化为任务结果;失败 → FAILED 并记录错误;
  • 恢复POST /api/v2/workflows/{job_id}/resume(与后端相同的请求/响应契约,含 404/409 语义以及针对 allowed_decisions 之外动作的 422 INVALID_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.events URL。它在 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}/events SSE 流回放已完成的运行,在挂起时呈现 human_input_request,并在恢复后带 Last-Event-ID 重连时只投递恢复后的 human_input_decision + end 帧。

6.8 组件索引(Concern → 代码位置)

关注点 位置 关键符号
暂停信号 + 检查点 lfx/graph/graph/base.pylfx/graph/checkpoint/ GraphPausedExceptionresume_from_checkpointset_run_id
构建缝隙 + 恢复重建 api/build.py build_resumed_graph_and_get_order_rerun_non_input_predecessors
恢复调度守卫 lfx/graph/graph/base.py is_vertex_runnablecheckpoint_restored_built_ids
Run-id 钉住 api/utils/flow_utils.py build_graph_from_data
持久化任务 + 单次恢复 services/background_execution/{runner,service}.py JobStatus.SUSPENDEDclaim_suspended_for_resume
lfx serve 持久化基座 lfx/services/durable/lfx/cli/serve_durable.py SqliteDurableJobStoreSqliteCheckpointStoreDurableServeWorkflowHost
追踪 services/tracing/{native,service}.py record_event_span_finalize_pending_spans_stop 排空
Trace UI TraceComponent/TraceDetailView.tsx backendHasGateexecutedOutputSpans
实时 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、无状态翻动。可安全上线;
  • Schemacheckpoint 表与 span/trace 表支撑持久化与可观测性;本特性不包含破坏性迁移;
  • 运行模式:在持久化后台路径(POST /api/v2/workflowsmode: 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]

十、延伸阅读与验证入口

想进一步在仓库中追踪本特性的实现与测试,建议按以下入口逐层深入:

适用前提与限制:HITL 暂停-恢复仅适用于持久化后台执行路径(mode: background)或显式开启 LFX_SERVE_DURABLE_DBlfx serve;嵌套 Flow 中已连线的 HITL 节点会被显式拒绝而非静默跳过;决策动作由用户自定义(不限于 approve/reject),超时后的迟到决策会被确定性地改道到 fallback 分支或哨兵动作,而不会被丢弃。

登录后查看全文
热门项目推荐
相关项目推荐