跳到主要内容

数据截至 (上游 commit b78a3462c9a6)

停下来等人:暂停、审批表单与恢复执行

30 秒导读: 工作流跑到「人工输入」节点会停住。Dify 把这一刻的引擎状态整个序列化存进对象存储, 给相关的人发一张表单(网页或邮件),然后放掉这次 HTTP 请求。等有人填了表单,一个 Celery 任务 把状态读回来、重建引擎、从断点继续跑完。本章讲这套「存档—叫人—读档」是怎么实现的。


1. 这是什么(零基础也能懂)

一句话定义: 让一次工作流运行可以中途停下来等真人拍板,并且这个「等」可以持续几天。

为什么需要它。 前面 01 章讲的那条链路是同步的:一个 HTTP 请求进来, 引擎跑完,SSE 流推完,请求结束。但真实业务里有些步骤机器不能自己决定——退款要不要批、这份文案能不能发、 这条 SQL 要不要真执行。这些地方需要一个人看一眼、点个按钮。

难在哪。 人不会在 30 秒内回你。一次审批可能拖到明天,甚至下周。这就意味着:

  • HTTP 连接不能一直挂着;
  • 跑到一半的引擎状态(变量池、已执行节点、token 计数)必须存下来,进程重启也不能丢;
  • 恢复时可能已经换了一台机器、换了一个进程,得能原样重建。

一句话直觉: 把它当成游戏存档。跑到人工节点时按下存档键,把整个内存快照写到磁盘; 玩家(审批人)什么时候回来都行,回来后从存档点继续,而不是从头再来一遍。

用起来什么样。 调用方(比如你自己的后端)通过 Service API 发起一次运行,SSE 流里会先收到常规的 node_started / node_finished,然后收到两个特殊事件:

event: human_input_required
data: {"form_id":"019...","node_id":"1712...","node_title":"经理审批",
"form_content":"请确认退款 ¥3200","form_token":"7Kd2f...","actions":[...],
"approval_channels":[],"expiration_time":1765000000}

event: workflow_paused
data: {"workflow_run_id":"018...","paused_nodes":["1712..."],
"reasons":[{"TYPE":"human_input_required","form_id":"019...","form_token":"7Kd2f..."}]}

流到这里就结束了。审批人拿着 form_token 去填表单并提交:

# 示意,非源码
curl -X POST "$BASE/v1/form/human_input/7Kd2f..." \
-H "Authorization: Bearer $APP_TOKEN" \
-d '{"action":"approve","inputs":{"comment":"同意"},"user":"u-42"}'

提交成功后工作流在后台继续跑。调用方可以重新连上事件流看后续结果(见 §3.6)。

本节不出现代码细节。往下才是原理。


2. 顶层全景(它大概怎么转)

怎么读这张图: 从左到右是时间顺序。上半行是「停下来」,下半行是「被叫醒」; 中间那条竖线是进程边界——左边是原来那次 HTTP 请求的线程,右边是几小时后另一个 Celery worker。

── 停下来 ───────────────────────┐ ┌────────── 被叫醒 ──────────
│ │
① 节点抛暂停 ② 存现场 │ │ ⑤ 有人提交表单
GraphRunPausedEvent → 序列化状态 │ │ submit_form_by_token
(reasons=[...]) 写对象存储 │ │ │
│ │ │ │ ↓
↓ ↓ │ │ ⑥ 丢进 Celery 队列
③ 建表单 + 投递 workflow_pauses │ │ resume_app_execution
human_input_forms (只存 object key) │ │ │
→ 网页链接 / 邮件 │ │ ↓
│ │ │ ⑦ 读回快照、重建引擎
↓ │ │ → generator.resume(...)
④ 推 SSE 后收流 │ │ │
请求结束 │ │ ↓
│ │ ⑧ 从断点继续跑完

部件一句话职责:

部件干什么在哪个文件
PauseReason / HumanInputRequired描述「为什么停」的类型外部包 graphon.entities.pause_reason
PauseStatePersistenceLayer听到暂停事件就把现场序列化落库api/core/app/layers/pause_state_persist_layer.py:77
WorkflowResumptionContext恢复所需的全部东西打成一个 JSONapi/core/app/layers/pause_state_persist_layer.py:42
WorkflowPause / WorkflowPauseReason暂停记录与暂停原因两张表api/models/workflow.py:2117 / :2100
HumanInputForm 及其投递/收件人表表单内容、发给谁、每人一个 tokenapi/models/human_input.py:28
HumanInputSurface + 允许表哪个 API 入口能操作哪类收件人api/core/workflow/human_input_policy.py:19
HumanInputService校验提交、落库、把恢复任务丢进队列api/services/human_input_service.py:157
_resume_app_executionCelery 侧的读档与重新起跑api/tasks/app_generate/workflow_execute_task.py:484
build_workflow_event_stream断线重连时把已发生的事件重讲一遍api/services/workflow_event_snapshot_service.py:74

一个必须先说清的边界。 graphon 是 Dify 通过 PyPI 引入的外部包api/pyproject.toml:48 锁在 graphon==0.7.0),克隆里没有它的源码。所以本章讲的是 Dify 侧怎么消费 这些类型—— PauseReasonHumanInputRequiredGraphRunPausedEventGraphRuntimeState内部实现读不到, 凡是涉及它们内部行为的地方本章都会明说。人工输入节点本身也住在 graphon 里 (graphon.nodes.human_input.*),Dify 只提供它需要的运行时能力,见 03 章


3. 核心原理(逐个机制)

3.1 暂停是怎么发生的

要解决的小问题: 引擎怎么把「我停了,而且是因为这个原因停的」告诉外面。

思路。 不用异常,用事件 + 原因列表。引擎发一个 GraphRunPausedEvent,里面带一个 reasons 序列;每个 reason 是一个 PauseReason 子类型。Dify 侧只做两件事:把 reason 落库、 按 reason 的类型决定要不要去叫人。

目前 Dify 认识两种 reason(依据:api/models/workflow.py:2232WorkflowPauseReason.from_entity 只对这两种做了 match,其余直接 raise AssertionError):

Reason 类型什么时候出现带的关键字段
HumanInputRequired人工输入节点建好表单、需要真人填form_idnode_idnode_titleform_contentinputsactions
SchedulingPause调度性暂停(例如时间片切换)message

Dify 侧的第一个消费点在 runner 的事件分发里:

case GraphRunPausedEvent():
runtime_state = workflow_entry.graph_engine.graph_runtime_state
paused_nodes = runtime_state.get_paused_nodes()
self._enqueue_human_input_notifications(event.reasons)

api/core/app/apps/workflow_app_runner.py:435)它做了两件事:把暂停事件转成队列事件 QueueWorkflowPausedEvent 往 SSE 管道推,以及调 _enqueue_human_input_notifications:698)——后者只挑出 HumanInputRequired,给每个 form_id 丢一个邮件投递任务 (dispatch_human_input_email_task.apply_async:705)。

落库映射。 reason 存进 workflow_pause_reasons 表时被拍平成三个字段 (api/models/workflow.py:2232 from_entity):type_ 记类型,form_id 只在人工输入时非空, message 只在调度暂停时非空。读回来时 to_entity():2161)重建对象,但注意 它只填得回 form_idnode_idform_content / node_title 一律给空串—— 真正的表单内容不在这张表里,在 human_input_forms.form_definition。仓储层的 _hydrate_pause_reasonsapi/repositories/sqlalchemy_api_workflow_run_repository.py:1081) 就是专门去补这一块的:按 form_id 批量捞表单和收件人,重建出完整的 HumanInputRequired

另一个消费点是策略层。 api/core/workflow/human_input_policy.py:161resolve_human_input_pause_reason_inputs 会拿变量池把 reason 里「选项来自变量」的下拉框 展开成真实选项值(resolve_variable_select_input_options:113)。注释里作者自己承认这是 a dirty hacks,但好处是调用方拿到的永远是具体选项而不是未解析的 selector。

3.2 现场怎么存下来

要解决的小问题: 恢复一次运行,到底需要存哪些东西?

思路。 需要两样,缺一不可:

  1. 引擎状态 GraphRuntimeState——变量池、已跑到哪、token 计数。这是「跑到哪了」。
  2. generate entity——发起这次运行的入参对象(app 配置、inputs、files、user_id、invoke_from…)。 这是「当初是怎么发起的」。

两样打成一个 WorkflowResumptionContextapi/core/app/layers/pause_state_persist_layer.py:42):

class WorkflowResumptionContext(BaseModel):
version: Literal["1"] = "1"
generate_entity: _GenerateEntityUnion
serialized_graph_runtime_state: str

引擎状态存的是字符串(graphon 自己 dumps() 出来的快照),Dify 不解释它的内容; generate entity 则是 pydantic 模型直接嵌进来。

为什么要包一层 wrapper。 只有 workflow 和 chatflow 两种 app 能暂停,对应两个不同的 entity 类。 反序列化时得知道该还原成哪个,所以外面套了带 type 判别字段的 wrapper (_WorkflowGenerateEntityWrapper / _AdvancedChatAppGenerateEntityWrapper:20/:25), 用 pydantic 的 Field(discriminator="type") 做判别联合(:30)。文件顶部的注释把这个动机写得很直白。

落库的执行者是一个 Layer。 PauseStatePersistenceLayer:64)挂在引擎的 Layer 链上 (Layer 机制见 04 章),on_event 里只关心一种事件:

if not isinstance(event, GraphRunPausedEvent):
return

:110)命中后组装 WorkflowResumptionContext,从变量池里取出 SystemVariableKey.WORKFLOW_EXECUTION_ID 作为 run id(:124),然后调 repo.create_workflow_pause(...):130)。

它自己不知道该用哪个 session、状态文件算谁的——这两个参数由 PauseStateLayerConfig:57)这个 frozen dataclass 从外面注入,字段只有两个: session_factorystate_owner_user_id(后者一般是工作流创建者,用于状态文件的归属)。 调用方在 _generate 里判断 pause_state_config is not None 才把这个 Layer 挂上去 (api/core/app/apps/workflow/app_generator.py:369)——所以没配置就等于关掉暂停持久化, 比如 workflow-as-tool 的嵌套调用就显式传了 pause_state_config=Noneapi/core/tools/workflow_as_tool/tool.py:130)。

真正写盘的地方在仓储层。 create_workflow_pauseapi/repositories/sqlalchemy_api_workflow_run_repository.py:967)做了这么几步:

校验 run 状态 ∈ {RUNNING, PAUSED}

删掉这个 run 之前的 pause 记录(一个 run 只允许一条)

state_obj_key = f"workflow-state-{uuid4()}.json"
storage.save(state_obj_key, state.encode()) ← 大 JSON 进对象存储

WorkflowPause(state_object_key=...) + N 条 WorkflowPauseReason

workflow_run.status = PAUSED

一个事务里全部提交

关键取舍:数据库里只存对象存储的 key,不存快照本身:993-994api/models/workflow.py:2256)。 变量池可能很大,塞进 Postgres 是灾难。

两张表的字段。

表 / 字段含义备注
workflow_pauses.workflow_id哪个版本的工作流注释解释了为什么不能只靠 app_id:一个 app 有多版本工作流,恢复时必须加载同一版
workflow_pauses.workflow_run_id哪次运行加了 UniqueConstraint,用唯一约束替代在大表 WorkflowRun 上加字段的迁移
workflow_pauses.state_object_key快照文件名真正内容在对象存储
workflow_pauses.resumed_at恢复时间恢复不删记录,只打时间戳
workflow_pause_reasons.type_ / form_id / message / node_id拍平后的 reasonform_id 仅当 type_ == HUMAN_INPUT_REQUIRED 时非空

一个埋伏笔。 AppGenerateEntity.trace_manager 被声明成 Field(default=None, exclude=True, repr=False)api/core/app/entities/app_invoke_entities.py:147)—— exclude=True 意味着它不会被序列化进快照。这个坑在 §3.5 结账。

顺便说一个观察:SuspendLayer 目前是死代码。 api/core/app/layers/suspend_layer.py:7 定义了一个 只维护布尔标志、暴露 is_paused():31)的 Layer,但全仓库除了它自己的单元测试没有任何引用grep -rn SuspendLayerapi/ 下只命中定义处和 api/tests/unit_tests/core/app/layers/test_suspend_layer.py)。 类的 docstring 是空的 """ """。真正干活的是 PauseStatePersistenceLayer

3.3 怎么把人叫来

要解决的小问题: 停下来之后,怎么知道该找谁、用什么方式找、他点进来看到什么。

三张表的关系。 一张表单可以有多种投递方式,每种方式下有多个收件人,每个收件人一个独立 token

HumanInputForm 一次暂停 = 一张表单
(表单内容/有效期/状态)
│ 1

│ N
HumanInputDelivery 一种投递方式(webapp / email)
(channel_payload = 该方式的完整配置 JSON)
│ 1

│ N
HumanInputFormRecipient 一个具体收件人
(recipient_type + access_token ← 每人一把钥匙)

api/models/human_input.py:28 / :98 / :221access_token:240, 默认值来自 _generate_token(),22 字符随机串,带 unique=True。)

每人一把钥匙的意义在两处兑现:邮件里给每个人拼的链接是 {APP_WEB_URL}/form/{token}api/tasks/mail_human_input_delivery_task.py:46), 以及提交时能记录 completed_by_recipient_id——谁批的可追溯。

投递方式只有两种。 DeliveryMethodTypeapi/core/workflow/human_input_adapter.py:27) 是 WEBAPPEMAIL。两者被建模成判别联合 DeliveryChannelConfig:176):

类型模型配置内容
webappInteractiveSurfaceDeliveryMethod:152空配置——网页表单不需要额外参数,链接本身就是投递
emailEmailDeliveryMethod:157EmailDeliveryConfig:69):收件人、主题、正文、debug_mode

邮件收件人这一族模型:41:66):

模型表示字段
BoundRecipient工作区内的成员(按 id 引用)reference_id
ExternalRecipient工作区外的人(按邮箱)email
EmailRecipients上面两种的容器include_bound_group(发给整个工作区)+ items

这里藏着一段兼容史。 include_bound_groupAliasChoices("include_bound_group", "whole_workspace") 同时接受老字段名(:62), 而 _normalize_email_recipients:344)在数据进 graphon 模型之前把老 payload 里的 whole_workspaceinclude_bound_groupuser_idreference_id 就地改名。整个 human_input_adapter.py 文件的开篇注释说明了这个模块存在的理由:画布和数据库里还躺着 Dify 自己的旧拼写, 在交给 graphon 之前统一翻译成新契约。这也是它叫 adapter 的原因。

邮件正文的渲染是三段流水线,顺序很重要(api/tasks/mail_human_input_delivery_task.py:103 _render_body):

原始 body(markdown,可含 {{#url#}} 和变量占位符)

│ ① replace_url_placeholder 把 {{#url#}} 换成该收件人专属的表单链接

│ ② variable_pool.convert_template 套用运行时变量(金额、申请人…)

│ ③ markdown → HTML → bleach 白名单清洗

可发送的 HTML

①② 在 EmailDeliveryConfig.render_body_templateapi/core/workflow/human_input_adapter.py:110), ③ 在 render_markdown_body:122)。第 ③ 步值得细看:bleach.clean(tags=[], strip=True) 把用户写的原始 HTML 剥干净, 渲染 markdown,最后 用白名单二次清洗——白名单只放行 16 个标签、a[href,title] 等三组属性,协议白名单在默认基础上多加一个 mailto:71:94)。

主题行单独处理:sanitize_subject:138)先把 \r \n 换成空格再清标签并压缩空白—— 这是防 邮件头注入(在主题里塞换行伪造额外邮件头)。

变量池从哪来? 邮件是异步发的,那时原进程早没了。所以 _load_variable_poolapi/tasks/mail_human_input_delivery_task.py:116)反过来去读刚存下的暂停快照: 读 WorkflowPauseWorkflowResumptionContext.loadsGraphRuntimeState.from_snapshot → 拿 variable_pool。 这份暂停快照在本章会被四处代码读回来:发邮件(这里)、渲染表单选项(HumanInputService._load_variable_pool_for_formapi/services/human_input_service.py:308)、恢复执行、重连重建事件流。

还有两类「隐式收件人」不来自画布配置,由仓储层在建表单时补上(api/core/repositories/human_input_repository.py:447 create_form):

隐式收件人何时创建判定函数
CONSOLE在 console 里调试(invoke_source == "debugger"),或在 explore 里跑且启用了 webapp 投递_should_create_console_recipient:419
BACKSTAGE只要是 RUNTIME 表单且有 invoke_source 或提交者身份_should_create_backstage_recipient:433

对应地,节点运行时在 console/explore 场景下会把 webapp 投递方式过滤掉api/core/workflow/node_runtime.py:857 _resolve_delivery_methods)——调试时不该真去生成对外的公开表单链接。

3.4 谁有权提交

要解决的小问题: 同一张表单有好几个 token,不同的 API 入口不能通用。

思路:按「入口」收窄「能操作的收件人类型」。 三个入口枚举在 HumanInputSurfaceapi/core/workflow/human_input_policy.py:19),允许关系是一张硬编码的表 (ALLOWED_RECIPIENT_TYPES_BY_SURFACE:23):

入口 surface认证方式允许操作的 recipient 类型潜台词
service_apiapp API tokenSTANDALONE_WEB_APPtoken 调用方只能碰面向终端用户的 web 表单
openapiapp API tokenSTANDALONE_WEB_APP同上
console登录态 + 租户CONSOLEBACKSTAGE内部流程只走 console 认证链

源码里那段注释把设计意图直说了:SERVICE_API and OPENAPI are intentionally narrower than CONSOLE。 落到 controller 上就是三行同构的守卫:

  • api/controllers/service_api/app/human_input_form.py:70 — 不匹配就抛 404(不是 403,不泄露表单是否存在
  • api/controllers/openapi/human_input_form.py:59
  • api/controllers/console/human_input_form.py:82

一张表单多个 token,对外该露哪一个? 用一张固定优先级表选(_RECIPIENT_TOKEN_PRIORITY:31, 数字小的优先):

优先级recipient 类型对应审批渠道 ApprovalChannel
0BACKSTAGEconsole
1CONSOLEconsole
2STANDALONE_WEB_APPweb_app
—(不参与)EMAIL_MEMBER / EMAIL_EXTERNALemail

get_preferred_form_token:47)就是遍历取最小优先级;注意邮件类收件人不在这张表里priority is None 时直接 continue——邮件 token 永远不会被 API 响应露出来,它只出现在发出去的邮件链接里。

拿不到 token 的人也要被告知去哪批。 这是 FormDisposition:62)的用途, 由 disposition_for_surface:74)算出来,两个字段:

该表单的全部 recipients

├── 本 surface 允许的 ──→ 取优先级最高的 token ──→ form_token

└── 本 surface 不允许的 ──→ 收集它们的 approval_channel(去重排序)──→ approval_channels

于是 Service API 调用方在 human_input_required 事件里可能拿到 form_token=nullapproval_channels=["console"]——明确告诉它「这单得去控制台批,别在这儿等」,而不是给一个语焉不详的空值。 enrich_human_input_pause_reasons:91)负责把这两个字段(外加 expiration_time)注回 reason 的 JSON 里。

从 surface 到具体请求的映射有两处:SSE 直出路径按 invoke_from 映射 (_INVOKE_FROM_TO_HITL_SURFACEapi/core/app/apps/common/workflow_response_converter.py:87), 重连路径由 controller 显式传(如 api/controllers/service_api/app/workflow_events.py:158)。 数据库查询集中在 load_form_dispositions_by_form_idapi/core/workflow/human_input_forms.py:24):按 form_id 批量捞收件人,逐表单算 disposition。

3.5 恢复路径

要解决的小问题: 表单提交之后,怎么把那次运行重新拉起来。

先看整条链。 竖线仍是进程边界:

HTTP 提交请求 │ Celery worker(另一个进程)
───────────── │ ─────────────────────────
submit_form_by_token │
│ ensure_form_active(已提交/超时/过期?)
│ _validate_submission(文件重建 + graphon 校验)
│ mark_submitted(写 submitted_data / status=SUBMITTED)
↓ │
enqueue_resume(workflow_run_id) ──────┼──→ _resume_app_execution(payload)
(apply_async 到队列) │ │ 读 WorkflowPause → 快照
│ │ GraphRuntimeState.from_snapshot
│ │ 捞 workflow / app / user
│ │ resume_workflow_pause(状态守卫)
│ ↓
│ generator.resume(...)
│ │ 重建 trace_manager
│ ↓
│ _generate(graph_runtime_state=...)
│ │ 引擎从断点继续
│ ↓
│ _publish_streaming_response → 事件打回 topic

入口:submit_form_by_tokenapi/services/human_input_service.py:195)。 它先按 token 取表单并核对 recipient_type 必须与调用方声称的一致,然后过三道闸:

  1. ensure_form_active:233)——已提交?状态是 TIMEOUT/EXPIRED?过了 expiration_time?超了全局超时?
  2. _validate_submission:248)——委托给 validate_and_normalize_submission:336)。 这里的分工写在 docstring 里:graphon 拥有表单 schema 与校验规则,Dify 负责租户感知的文件重建。 即先由 Dify 把提交里的文件 mapping 用 build_from_mapping 重建成受租户校验的 File 对象(:427), 再交给 graphon_validate_human_input_submission 校验形状。
  3. mark_submittedapi/core/repositories/human_input_repository.py:611)——写入 selected_action_id / submitted_data / status=SUBMITTED / completed_by_recipient_id

表单是一次性的。 服务层的错误定义把语义写在描述里: This form has already been submitted by another userapi/services/human_input_service.py:109), HTTP 412。先到先得,无论是谁提交的。

分路由。 提交完成后按归属决定叫醒谁(:224:231):有 workflow_run_idenqueue_resume(工作流/聊天流),只有 conversation_idenqueue_agent_app_resume (Agent v2 聊天的 ask_human,另一条链路,本章不展开)。enqueue_resume:265) 还会核对 app mode 属于 {WORKFLOW, ADVANCED_CHAT},否则只打一条 warning。

Celery 侧的读档_resume_app_executionapi/tasks/app_generate/workflow_execute_task.py:484):

resumption_context = WorkflowResumptionContext.loads(pause_entity.get_state().decode())
generate_entity = resumption_context.get_generate_entity()
graph_runtime_state = GraphRuntimeState.from_snapshot(resumption_context.serialized_graph_runtime_state)

:481:488)后面把 workflow / app / user 从库里重新捞出来;chatflow 还要额外捞回 Conversation 和这次 run 对应的最后一条 Message:518:533),任何一个缺失都只是 warning + return。

先翻状态,再起跑。 resume_workflow_pause:546,实现在 api/repositories/sqlalchemy_api_workflow_run_repository.py:1156)在一个事务里做了四重守卫, 任何一条不满足就抛错:

检查防的是什么
workflow_run.status == PAUSED重复恢复一个已经在跑的 run
workflow_run.pause is not None状态与记录不一致
pause_model.id == pause_entity.id拿着旧的 pause 实体去恢复新的暂停
pause_model.resumed_at is None同一个暂停被恢复两次

通过后把 resumed_at 打上时间戳、workflow_run.status = RUNNING:1176:1177)。

两个 resume 入口是对称的——api/core/app/apps/workflow/app_generator.py:274api/core/app/apps/advanced_chat/app_generator.py:272,签名都收 graph_runtime_statepause_state_config,最后都转调各自的 _generate(...) 并把 graph_runtime_state 透传下去 (api/core/app/apps/workflow/app_generator.py:318)。「恢复」不是一条特殊代码路径, 只是「首跑」多传了一个已有状态——首跑时这个参数是 None

为什么必须重建 trace_manager。 两个 resume 里都有同一段代码和同一段 docstring:

if application_generate_entity.trace_manager is None:
application_generate_entity = application_generate_entity.model_copy(
update={"trace_manager": TraceQueueManager(app_id=..., user_id=...)}
)

api/core/app/apps/workflow/app_generator.py:295;chatflow 同构在 :278) 原因就是 §3.2 埋的伏笔:trace_manager 声明为 exclude=Trueapi/core/app/entities/app_invoke_entities.py:147),它是个持有队列和线程的活对象,序列化不进去也不该进去。 从快照恢复出来的 entity 这个字段必然是 None,不补上的话恢复后的执行就没有链路追踪 (追踪机制见 04 章)。docstring 写得很清楚: trace_manager is transient and excluded from generate-entity serialization

注意补的时机——在把 entity 交给持久化 Layer 之前。因为下一次暂停时 PauseStatePersistenceLayer 还要把这个 entity 再存一遍,但 exclude=True 保证它又会被剔掉。 一次运行可以暂停多次,这个环闭得上。

恢复后的流去哪了。 恢复时没有 HTTP 客户端在等,但 resumed_generate_entity 被强制 stream=Trueapi/tasks/app_generate/workflow_execute_task.py:628)。产出的 generator 交给 _publish_streaming_response:342),逐个事件 topic.publish 到 Redis 主题上, 供后来重连的客户端订阅。它的 docstring 点明了它的兜底职责:一旦开始迭代运行时流, 这里就是最后能保证 SSE 消费者看到终结事件的地方——流异常中断时它会补发一个 FAILED 的 workflow_finished:359 _publish_failed_terminal_event)。终结事件集合是 {"workflow_finished", "workflow_paused"}:403):再次暂停也算一种正常终结。

一处不对称。 workflow 分支跑完后会 delete_workflow_pause 清掉暂停记录和状态文件 (:712,并容忍「记录已被替换或删除」的情况),chatflow 分支没有这一步_resume_advanced_chat:586:646)。代码里没有解释原因。

3.6 断线重连:服务端把已发生的事重讲一遍

要解决的小问题: 客户端断开过、或者干脆是几小时后新连上来的,它错过的事件怎么补。

思路:不重放日志,而是从当前落库状态「重新编」一份等价的事件序列。 build_workflow_event_streamapi/services/workflow_event_snapshot_service.py:74) 先订阅实时 topic 并开一个后台线程缓冲_start_buffering:552,避免边补历史边丢新事件), 然后由 _build_snapshot_events:220)按顺序造出:

workflow_started ← 从 WorkflowRun 造

[message_replace] ← 仅 chatflow,把已有 answer 补上

node_started / node_finished ← 逐条节点执行快照;仍 RUNNING 的只发 started

human_input_required × N ← 仅当 run 处于 PAUSED

workflow_paused ← 仅当 run 处于 PAUSED

(切到实时缓冲队列,继续转发)

契约必须一模一样。 _build_pause_event:470)和实时路径的 workflow_pause_to_stream_responseapi/core/app/apps/common/workflow_response_converter.py:319) 调用的是同一组 helper:resolve_human_input_pause_reason_inputsload_form_dispositions_by_form_idenrich_human_input_pause_reasons。两处都写着同一句注释: Reconnect paths must preserve the same pause-reason contract as live streams; otherwise clients see schema drift after resume. 这就是这三个函数被抽到 human_input_policy.py / human_input_forms.py 而不是内联的原因。

task_id 的连续性也要保。 客户端可能拿 task_id 做关联,所以 _resolve_task_id:201优先从暂停快照里的 generate entity 取原来的 task_id,取不到才退到缓冲里的提示值, 再不行才用 run id 顶替。这是那份暂停快照被读回来的又一处。

行为开关。 close_on_pause 决定发完 workflow_paused 就收流还是继续挂着,由 query 参数 continue_on_pause 反向控制(api/controllers/service_api/app/workflow_events.py:146); include_state_snapshot=true 才走快照重建,否则走普通的实时事件重取。


4. 服务层与异步执行(周边设施)

主链路之外还有几个服务把这套机制撑起来:

文件干什么关键符号
api/services/human_input_service.py表单读取、校验、提交、排队恢复HumanInputService:157submit_form_by_token:195
api/services/human_input_file_upload_service.py表单里传文件用的独立令牌体系issue_upload_token:72validate_upload_token:102
api/services/human_input_delivery_test_service.py编辑器里「发一封测试邮件」HumanInputDeliveryTestService:107EmailDeliveryTestHandler:120
api/services/workflow_event_snapshot_service.py重连时重建事件流build_workflow_event_stream:69
api/services/async_workflow_service.py让运行本身就异步起跑(按订阅档位分队列)trigger_workflow_async:52
api/tasks/human_input_timeout_tasks.py定时扫过期表单check_and_handle_human_input_timeouts:56

上传令牌为什么要另起一套。 HumanInputFormUploadToken 的 docstring 说得很直接: HITL 上传令牌故意与 app/service bearer token 分开,存成不透明随机值,让上传端点做一次直接查表 而不进入正常的 Web App 认证链api/models/human_input.py:283)。令牌本身不存归属, 归属在校验时从表单所属的工作流发起人反推(_resolve_upload_ownerapi/services/human_input_file_upload_service.py:149)——这样恢复时重建文件仍然走正常的文件访问检查。 投递测试表单没有工作流运行,其上传归到 app 创建者名下(_resolve_delivery_test_upload_owner:198)。

投递测试是编辑期能力:EmailDeliveryTestHandler.send_testapi/services/human_input_delivery_test_service.py:132)复用了和真投递完全相同的三段渲染 (render_body_templaterender_email_templaterender_markdown_body),所以「测试时好看、真发时变形」这类问题基本被堵死。它另外做了两道前置检查:租户套餐是否开了邮件投递、mail client 是否初始化。

超时是两级的。

级别判定后果
节点级HumanInputForm.expiration_time 到了表单标 TIMEOUT照样恢复工作流(把超时作为结果送回引擎)
全局级created_at + HUMAN_INPUT_GLOBAL_TIMEOUT_SECONDS 到了表单标 EXPIRED直接终结整个 run

全局超时的处理在 _handle_global_timeoutapi/tasks/human_input_timeout_tasks.py:32): 把 WorkflowRun 置为 STOPPED、写 error、删掉对象存储里的状态文件、给 pause 打上 resumed_at (相当于作废)。默认全局上限是 7 天api/configs/feature/__init__.py:98), 扫描任务由 Celery beat 定期触发(api/extensions/ext_celery.py:235)。

async_workflow_service 与本章的关系是间接的:它让一次运行从一开始就在 Celery 里跑 (trigger_workflow_asyncapi/services/async_workflow_service.py:55,按订阅档位分派到 professional / team / sandbox 三个队列)。这类运行本来就没有 HTTP 客户端在等, 所以它和恢复执行走的是同一种「事件发布到 topic、客户端另行订阅」的形态。 触发器那一半入口见 06 章


5. 巧妙之处(可以带走的技术)

  • 「恢复」不是特殊路径,只是多传一个状态。 resume() 和首跑最终汇进同一个 _generate(...), 差别只在 graph_runtime_state 是否为 Noneapi/core/app/apps/workflow/app_generator.py:274 vs :312)。 想让一个执行引擎支持断点续跑,这是成本最低的接法。

  • 大状态进对象存储,库里只留一把钥匙。 state_object_key + storage.saveapi/repositories/sqlalchemy_api_workflow_run_repository.py:1015),变量池再大也不压 Postgres。

  • exclude=True 显式标注「不可序列化的活对象」,并在恢复处集中补齐。 trace_managerapi/core/app/entities/app_invoke_entities.py:147)是教科书式的例子: 声明处一行字段配置,恢复处一段带解释的重建,中间不需要任何人记住这条规矩。

  • 权限收窄写成一张数据表而不是散在 if 里。 ALLOWED_RECIPIENT_TYPES_BY_SURFACEapi/core/workflow/human_input_policy.py:27)让「service API 只能碰 web 表单」这条策略 可以一眼读完、一处修改,三个 controller 只是同构地调用同一个判定。

  • 拒绝时给出替代路径而不是空值。 FormDisposition.approval_channels:62)把 「你不能在这批,但可以去 console 批」变成了结构化字段。

  • 实时流与重连流复用同一组组装函数。 两处相同的注释明确了这是为了防止 schema drift (api/services/workflow_event_snapshot_service.py:601api/core/app/apps/common/workflow_response_converter.py:368)。

  • 邮件渲染的清洗顺序:剥 → 渲染 → 白名单。bleach.clean(tags=[]) 剥掉用户原始 HTML, 再渲染 markdown,最后按白名单二次清洗(api/core/workflow/human_input_adapter.py:123), 外加主题行的换行剥离防头注入(:138)。


6. 边界与局限(诚实版)

  • Web 表单端点是刻意不认证的。 api/controllers/web/human_input_form.py:156 直写 this endpoint is unauthenticated on purpose for now。安全性完全建立在 token 不可猜 (22 字符随机、唯一)加 IP 级限流(_FORM_ACCESS_RATE_LIMITER / _FORM_SUBMIT_RATE_LIMITER:63:71)之上。 同一文件 :165 还留着一条 TODO:目前没有阻止用「只应在 console 使用」的 token 从 web 端点提交

  • 表单一次性、先到先得。 多个收件人拿到的是各自的 token,但表单只有一份状态;第一个提交的人生效, 其余收到 412(api/services/human_input_service.py:107)。代码里没有会签/多人审批语义。

  • 暂停原因只认两种。 WorkflowPauseReason.from_entity 遇到第三种直接 AssertionErrorapi/models/workflow.py:2256),扩展需要同时改模型和迁移。

  • WorkflowPauseReason.to_entity() 是有损的。 它填不回 form_content / node_title:2161), 必须靠 _hydrate_pause_reasons 回表补全;直接用 to_entity() 的调用方会拿到空串。

  • SuspendLayer 是死代码。 定义存在、测试存在,但生产路径无人引用 (api/core/app/layers/suspend_layer.py:7)。

  • chatflow 恢复后不清理 pause 记录,workflow 会(api/tasks/app_generate/workflow_execute_task.py:748)。 代码里未说明原因,可能是遗漏也可能是有意为之——看不出来

  • graphon 内部不可见。 人工输入节点、GraphRuntimeState.dumps() / from_snapshot() 的具体格式、 GraphRunPausedEvent 的触发条件,都在 graphon==0.7.0api/pyproject.toml:48)里, 本章无法基于源码断言。快照格式的版本号倒是 Dify 自己管的:WorkflowResumptionContext.version 当前固定 "1"api/core/app/layers/pause_state_persist_layer.py:42),但仓库里没有任何 处理 version != "1" 的迁移代码——跨版本恢复旧快照会怎样,代码里看不出来


7. 与本组其它章的关系

你想知道去哪章
一次普通运行怎么从 HTTP 走到 SSE01 请求生命周期
人工输入节点在画布上长什么样、DSL 里怎么存02 画布 / 图存储 / DSL
graphon 边界具体划在哪、节点怎么被造出来03 节点组装与 graphon 边界
Layer 钩子机制本身、普通执行期怎么落库留痕04 执行期横切与留痕
谁把这些运行触发起来、外部能力怎么接06 触发器与插件运行时

8. 代码地图(导航索引)

主题文件路径符号名
暂停快照容器api/core/app/layers/pause_state_persist_layer.pyWorkflowResumptionContext
暂停持久化 Layerapi/core/app/layers/pause_state_persist_layer.pyPauseStatePersistenceLayerPauseStateLayerConfig
未接线的暂停标志 Layerapi/core/app/layers/suspend_layer.pySuspendLayer
暂停记录与原因表api/models/workflow.pyWorkflowPauseWorkflowPauseReasonfrom_entityto_entity
暂停落库 / 恢复 / 删除api/repositories/sqlalchemy_api_workflow_run_repository.pycreate_workflow_pauseresume_workflow_pausedelete_workflow_pause_hydrate_pause_reasons
引擎暂停事件的消费api/core/app/apps/workflow_app_runner.py_enqueue_human_input_notifications
表单 / 投递 / 收件人模型api/models/human_input.pyHumanInputFormHumanInputDeliveryHumanInputFormRecipientRecipientTypeApprovalChannel
上传令牌与文件关联api/models/human_input.pyHumanInputFormUploadTokenHumanInputFormUploadFile
投递方式与邮件渲染api/core/workflow/human_input_adapter.pyDeliveryMethodTypeEmailDeliveryConfigBoundRecipientExternalRecipientEmailRecipients
旧字段兼容适配api/core/workflow/human_input_adapter.pyadapt_human_input_node_data_for_graph_normalize_email_recipients
surface 权限策略api/core/workflow/human_input_policy.pyHumanInputSurfaceALLOWED_RECIPIENT_TYPES_BY_SURFACEdisposition_for_surfaceFormDisposition
token 优先级与事件补字段api/core/workflow/human_input_policy.pyget_preferred_form_tokenenrich_human_input_pause_reasons
表单 disposition 查询api/core/workflow/human_input_forms.pyload_form_dispositions_by_form_id
表单创建(含隐式收件人)api/core/repositories/human_input_repository.pyHumanInputFormRepositoryImpl.create_form_should_create_console_recipient
提交 / 超时状态机api/core/repositories/human_input_repository.pymark_submittedmark_timeout
提交校验与排队恢复api/services/human_input_service.pysubmit_form_by_tokenensure_form_activeenqueue_resume
Celery 侧读档重跑api/tasks/app_generate/workflow_execute_task.py_resume_app_execution_resume_workflow_publish_streaming_response
两个 resume 入口api/core/app/apps/workflow/app_generator.pyapi/core/app/apps/advanced_chat/app_generator.pyWorkflowAppGenerator.resumeAdvancedChatAppGenerator.resume
邮件投递任务api/tasks/mail_human_input_delivery_task.pydispatch_human_input_email_task_render_body_load_variable_pool
超时扫描任务api/tasks/human_input_timeout_tasks.pycheck_and_handle_human_input_timeouts_handle_global_timeout
重连事件重建api/services/workflow_event_snapshot_service.pybuild_workflow_event_stream_build_snapshot_events_build_pause_event
实时暂停事件转换api/core/app/apps/common/workflow_response_converter.pyworkflow_pause_to_stream_response_INVOKE_FROM_TO_HITL_SURFACE
HITL 上传令牌服务api/services/human_input_file_upload_service.pyHumanInputFileUploadService.issue_upload_token
投递测试服务api/services/human_input_delivery_test_service.pyEmailDeliveryTestHandler.send_test
异步起跑(分队列)api/services/async_workflow_service.pyAsyncWorkflowService.trigger_workflow_async
四个 API 入口api/controllers/{web,console,service_api/app,openapi}/human_input_form.pyHumanInputFormApiConsoleHumanInputFormApiWorkflowHumanInputFormApi