数据截至 (上游 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 | 恢复所需的全部东西打成一个 JSON | api/core/app/layers/pause_state_persist_layer.py:42 |
WorkflowPause / WorkflowPauseReason | 暂停记录与暂停原因两张表 | api/models/workflow.py:2117 / :2100 |
HumanInputForm 及其投递/收件人表 | 表单内容、发给谁、每人一个 token | api/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_execution | Celery 侧的读档与重新起跑 | 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 侧怎么消费 这些类型——
PauseReason、HumanInputRequired、GraphRunPausedEvent、GraphRuntimeState 的内部实现读不到,
凡是涉及它们内部行为的地方本章都会明说。人工输入节点本身也住在 graphon 里
(graphon.nodes.human_input.*),Dify 只提供它需要的运行时能力,见 03 章。
3. 核心原理(逐个机制)
3.1 暂停是怎么发生的
要解决的小问题: 引擎怎么把「我停了,而且是因为这个原因停的」告诉外面。
思路。 不用异常,用事件 + 原因列表。引擎发一个 GraphRunPausedEvent,里面带一个
reasons 序列;每个 reason 是一个 PauseReason 子类型。Dify 侧只做两件事:把 reason 落库、
按 reason 的类型决定要不要去叫人。
目前 Dify 认识两种 reason(依据:api/models/workflow.py:2232 的 WorkflowPauseReason.from_entity
只对这两种做了 match,其余直接 raise AssertionError):
| Reason 类型 | 什么时候出现 | 带的关键字段 |
|---|---|---|
HumanInputRequired | 人工输入节点建好表单、需要真人填 | form_id、node_id、node_title、form_content、inputs、actions |
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_id 和 node_id,form_content / node_title 一律给空串——
真正的表单内容不在这张表里,在 human_input_forms.form_definition