跳到主要内容

数据截至 (上游 commit b78a3462c9a6)

入口的另一半与外部能力:触发器与插件运行时

30 秒导读: 前面几章讲的都是"人在页面上点了运行,然后发生了什么"。这一章补两块拼图—— 谁还能启动这张图(webhook / 定时 / 第三方事件),以及图里的节点凭什么真的会干活 (工具、模型、数据源的实现全在另一个进程里,Dify 只是它的 HTTP 客户端)。

本章的 path:line 引用都相对克隆根 dify/,后端 Python 代码都在 api/ 目录下。

想读相邻内容:一次运行从 HTTP 到 SSE 见 01-request-lifecycle.md; JSON 图怎么变成可执行图、节点对象怎么构造见 03-node-assembly-graphon-boundary.md; Layer 钩子与留痕落库的通用机制见 04-execution-layers-and-tracing.md


1. 这章解决什么问题(零基础也能懂)

先说结论:一张 Dify 工作流图,有两个外界依赖,而它们都不在图本身里。

第一个依赖:谁按下了开始。

默认答案是"用户在页面上点了运行"。但真实生产里更常见的是这三种:

启动方式白话场景对应节点类型
Webhook别的系统往一个 URL 上 POST 一坨 JSON,工作流就跑一次trigger-webhook
定时"每天早上 9 点跑一次日报"trigger-schedule
插件事件"Slack 有人 @我" / "GitHub 开了新 issue" 就跑一次trigger-plugin

第二个依赖:节点凭什么会干活。

一个"调用 Serper 搜索"的工具节点、一个"调用 Claude"的 LLM 节点,它们的实现代码不在 Dify 仓库里。 Dify 主进程只负责把参数拼好,通过 HTTP 发给一个叫 plugin daemon(插件守护进程) 的独立服务, 再把返回的流转成节点输出。

一句话直觉:把 Dify 主进程想成一个"排程中心 + 电话总机"——排程中心接外面打进来的电话(触发器), 总机把每个节点的活儿转接给真正干活的分机(插件守护进程)。


2. 顶层全景

三条外部入口,最后都汇进同一个口子。先看这张图,它是本章前半段的骨架:

怎么读: 左边三列是三条互不相干的入口路线,中间那道横线是它们的唯一汇流点, 右边是异步执行。注意左边两条走的是"HTTP 进来 → 落 Celery 队列",只有定时那条是"Celery beat 主动醒来"。

外部世界 Dify api 进程 Celery worker
───────────── ────────────────────────── ──────────────────────

别的系统 POST ──HTTP──> /triggers/webhook/<id>
WebhookService ──> trigger_workflow_async
解析+校验 body/header/query (webhook 队列)

第三方 SaaS ──HTTP──> /triggers/plugin/<endpoint_id>
(Slack/GitHub) TriggerService.process_endpoint
①问插件"这请求算哪些事件" ──> dispatch_triggered_
②立刻回 200,剩下的丢队列 workflows_async
③再问插件"把请求翻成变量"
④trigger_workflow_async

(无) ─beat 定时─> poll_workflow_schedules
扫 next_run_at <= now ──> run_schedule_trigger
trigger_workflow_async

╔══════════════════════════════════════════════════╗
║ AsyncWorkflowService.trigger_workflow_async ║
║ 写 WorkflowTriggerLog(PENDING) + 投递执行队列 ║
╚══════════════════════════════════════════════════╝

v
WorkflowAppGenerator.generate(root_node_id=触发节点id,
layers=[TriggerPostLayer])

部件一句话职责:

部件干什么在哪
三个触发节点类图上的可视化起点,运行时只是"把池子里的值摊开给下游"api/core/workflow/nodes/trigger_*/
TriggerManager和插件守护进程谈订阅:列 provider、订阅、退订、刷新、调事件api/core/trigger/trigger_manager.py:30
services/trigger/*把一次外部 HTTP / 一次定时到点,落成一次工作流运行api/services/trigger/
TriggerPostLayer运行结束后回填触发日志(状态、耗时、token、输出)api/core/app/layers/trigger_post_layer.py:24
core/plugin/impl/*对插件守护进程的一组 HTTP 客户端(工具/模型/数据源/触发器/Agent/端点/OAuth)api/core/plugin/impl/
backwards_invocation/*反方向:插件回头调 Dify(用模型、用工具、跑节点、调 app)api/core/plugin/backwards_invocation/

3. 三类触发节点:图的另一种起点

3.1 一个常量决定了"谁能当起点"

图的起点不是随便定的。graphon 侧只认几种 execution_type = ROOT 的节点, 而 Dify 侧用一个 frozenset 声明"这些类型可以当入口":

TRIGGER_NODE_TYPES: Final[frozenset[str]] = frozenset(
(
TRIGGER_WEBHOOK_NODE_TYPE,
TRIGGER_SCHEDULE_NODE_TYPE,
TRIGGER_PLUGIN_NODE_TYPE,
)
)

这是 api/core/trigger/constants.py:7TRIGGER_NODE_TYPES,三个字符串常量分别是 trigger-webhook / trigger-schedule / trigger-plugin(同文件 :3-5)。

它被 api/core/workflow/node_factory.py:23 导入,然后直接摊进起点集合

_START_NODE_TYPES: frozenset[NodeType] = frozenset(
(BuiltinNodeTypes.START, BuiltinNodeTypes.DATASOURCE, *TRIGGER_NODE_TYPES)
)

node_factory.py:81-83。判定函数 is_start_node_typenode_factory.py:167get_default_root_node_id 紧随其后用它挑默认入口。

要点: 加一种触发器 = 往 TRIGGER_NODE_TYPES 里加一个字符串,起点集合自动跟着变—— 不需要改工厂逻辑。另一个消费者 is_trigger_node_typeconstants.py:16)被 api/services/workflow_service.py:727api/services/workflow_draft_variable_service.py:1176 用来 判断"这节点的草稿变量该不该按起点方式处理"。

3.2 三个节点的 _run 都薄得出奇

初读会觉得反直觉:触发节点的 _run 里没有任何"接收外部请求"的代码。

因为真正的接收工作在图跑起来之前就做完了,值已经被塞进变量池。节点要做的只是把 自己前缀下的变量摊开成 outputs,供下游引用。看 TriggerScheduleNode._runapi/core/workflow/nodes/trigger_schedule/trigger_schedule_node.py:36):

node_inputs = dict(self.graph_runtime_state.variable_pool.get_by_prefix(self.id))
system_inputs = self.graph_runtime_state.variable_pool.get_by_prefix(SYSTEM_VARIABLE_NODE_ID)

TriggerEventNode._runtrigger_plugin/trigger_event_node.py:43)几乎一模一样,只多挂了一份 WorkflowNodeExecutionMetadataKey.TRIGGER_INFO 元数据(provider_id / event_name / plugin_unique_identifier),供留痕使用。

TriggerWebhookNode 是三个里唯一有真逻辑的(trigger_webhook/node.py:22)——因为它要按节点配置 挑字段_extract_configured_outputs(同文件 :106)遍历 node_data.headers / params / body, 从 webhook_data 里逐个取值,并把原始数据一并放进 outputs["_webhook_raw"]。 其中 header 名做了"横杠/下划线互查 + 小写兜底"的容错查找,body 里声明为 SegmentType.FILE 的参数 走 generate_file_var:81)转成 FileVariable

3.3 三份 entities 定义了"画布上能配什么"

节点entities 文件关键模型能配的东西
webhooktrigger_webhook/entities.pyWebhookData:85method、content_type、headers/params/body 三组参数声明、status_code、response_body、timeout
scheduletrigger_schedule/entities.pyTriggerScheduleNodeData:10VisualConfig:35mode(visual/cron)、frequency、cron_expression、timezone
plugintrigger_plugin/entities.pyTriggerEventNodeData:14plugin_id、provider_id、event_name、subscription_id、event_parameters

三份 entities 各有一处值得留意的约束:

  • webhook 对参数类型做了白名单。 header 只准 STRING,query 只准 STRING/NUMBER/BOOLEAN, body 才放开到 object/array/file(trigger_webhook/entities.py:11-33 三个 frozenset + 三个 field_validator)。这样才能保证从 URL/头里解析出来的字符串一定能安全转型。
  • schedule 的可视化配置最终一律编译成 cron。 VisualConfig 里的 on_minute / time / weekdays / monthly_daysScheduleService.visual_to_cronapi/services/trigger/schedule_service.py:257)翻成 cron 表达式,运行期只认 cron。
  • 插件触发节点的事件参数只支持常量。 TriggerEventNodeData.resolve_parameterstrigger_plugin/entities.py:53)在遇到 type != "constant" 时直接抛 TriggerEventParameterError——因为这些参数要在图还没跑起来时就用于调用插件, 那时变量引用还没有值可解。

4. 外部事件怎么落到一次运行

4.1 Webhook:最短的一条路

POST /triggers/webhook/<webhook_id> 打进 api/controllers/trigger/webhook.py:59handle_webhook

分四步,每步一个方法,都在 api/services/trigger/webhook_service.pyWebhookService:74)里:

  1. 找到人。 get_webhook_trigger_and_workflow:88)按 webhook_idWorkflowWebhookTrigger,再查 AppTrigger 的状态;状态是 RATE_LIMITED 或非 ENABLED 就直接 报错,最后取已发布工作流(debug 模式则取草稿版)。
  2. 解析并校验。 extract_and_validate_webhook_data:171)按 content-type 分派到 _extract_json_body / _extract_form_body / _extract_multipart_body / _extract_octet_stream_body / _extract_text_body,再按节点声明的类型逐个转换与校验。
  3. 拼 inputs。 build_workflow_inputs:774)返回四个键: webhook_data / webhook_headers / webhook_query_params / webhook_body
  4. 投递。 trigger_workflow_execution:791)先扣配额,再调 AsyncWorkflowService.trigger_workflow_async

第 4 步里有一段值得单独看的配额处理:

# 删掉缩进和 session 管理后的骨架,逐字源码见 webhook_service.py:823-846
quota_charge = QuotaService.reserve(QuotaType.TRIGGER, webhook_trigger.tenant_id)
try:
AsyncWorkflowService.trigger_workflow_async(session, end_user, trigger_data)
quota_charge.commit()
except Exception:
quota_charge.refund()
raise

预留 → 成功提交 / 失败退还。若 QuotaExceededError,先调 AppTriggerService.mark_tenant_triggers_rate_limitedapi/services/trigger/app_trigger_service.py:24) 把该租户所有 ENABLED 触发器批量改成 RATE_LIMITED,再抛出——下次请求在第 1 步就会被挡掉, 不用每次都跑到扣配额才失败。

调试端点 /triggers/webhook-debug/<id>controllers/trigger/webhook.py:95)走另一条路: 不落 Celery、不跑已发布工作流,只往 TriggerDebugEventBus 推一个内存事件给正在监听的变量检查器; 没有监听者时会显式报错而不是假装 200。

4.2 插件事件:两段式,中间隔一个队列

这条路最绕,因为Dify 不认识第三方的请求格式——它得问插件两次。

第三方 POST /triggers/plugin/<endpoint_id>

│ [api 进程,同步,要尽快返回]

TriggerService.process_endpoint
├─ 按 endpoint_id 查 TriggerSubscription
├─ 解密该订阅的 credentials
├─ ① controller.dispatch(...) ──HTTP──> 插件:"这个请求算哪些事件?"
│ 返回 events:[...] + payload + 要原样回给第三方的 response
├─ 把原始请求和 payload 存进对象存储(request_id 索引)
├─ 校验每个 event_name 在 provider 声明里存在
├─ dispatch_triggered_workflows_async.delay(...) ← 丢队列,立刻返回
└─ return dispatch_response.response ← 第三方拿到的响应

│ [Celery worker,异步]

dispatch_triggered_workflow(每个 event 一次)
├─ 取回存储里的 request + payload
├─ 查有哪些工作流订阅了这个 (subscription, event_name)
├─ ② TriggerManager.invoke_trigger_event(...) ──HTTP──> 插件:"翻成变量"
│ 返回 variables(就是 trigger 节点的 inputs)或 cancelled=True
└─ AsyncWorkflowService.trigger_workflow_async(...)

对应代码:

步骤位置
路由 + 处理链api/controllers/trigger/trigger.py:18 trigger_endpoint
同步段api/services/trigger/trigger_service.py:77 TriggerService.process_endpoint
请求存盘api/services/trigger/trigger_request_service.py:47 persist_request / :58 persist_payload
异步段(单事件)api/tasks/trigger_processing_tasks.py:232 dispatch_triggered_workflow
异步段(Celery 入口)api/tasks/trigger_processing_tasks.py:446 dispatch_triggered_workflows_async
找订阅方api/services/trigger/trigger_subscription_operator_service.py:11 get_subscriber_triggers

为什么要"存盘 + 队列"而不是一口气做完? 因为第三方 webhook 通常有几秒的响应超时,而一个 endpoint 可能同时喂给同租户下多个 app 的多张图;每张图都要单独调一次插件把请求翻成变量。 所以同步段只做"认领事件 + 回响应",把重活丢给队列。

原始请求怎么跨进程传?序列化成字节丢进对象存储:TriggerHttpRequestCachingServicetrigger_request_service.py:11)用 serialize_request / deserialize_requesttriggers/<request_id>.raw.payload 两个键上存取——队列里只传一个 request_id 字符串

两个容易漏掉的分支:

  • 插件可以主动说"这次不算"。 TriggerInvokeEventResponse.cancelled=True 时(比如插件自己判断 这条 GitHub 事件不该触发),trigger_processing_tasks.py:363-371 会退还配额并 continue, 不产生任何运行。TriggerManager.invoke_trigger_eventapi/core/trigger/trigger_manager.py:150) 还把插件抛的 EventIgnoreError 统一翻译成 cancelled=True:192-193)。
  • 插件调用失败要留痕。 PluginInvokeError 时走 _record_trigger_failure_logtrigger_processing_tasks.py:126),补写一条 workflow run + trigger log,用户在日志页能看见 这次失败,而不是事件凭空消失。

顺带一提,同一个同步段还会把事件推给调试总线:dispatch_trigger_debug_eventtrigger_processing_tasks.py:59)——这就是画布上"插件触发器调试"面板能实时看到事件的原因。

4.3 定时:两个 Celery beat 任务

定时侧没有外部 HTTP,全靠 beat 周期性醒来。注册在 api/extensions/ext_celery.py:266-277, 两个任务都由开关控制:

beat 任务干什么周期
poll_workflow_schedules扫到期的计划并派发执行WORKFLOW_SCHEDULE_POLLER_INTERVAL 分钟
trigger_provider_refresh扫快过期的订阅/凭据并派发刷新TRIGGER_PROVIDER_REFRESH_INTERVAL 分钟

排程轮询api/schedule/workflow_schedule_task.py:18)是一个批处理循环:

  • _fetch_due_schedules:55)用 WorkflowSchedulePlan JOIN AppTriggernext_run_at <= now 且触发器 ENABLED 的记录,按最超期优先排序, 关键是 .with_for_update(skip_locked=True):82)——多个 poller 并行也不会重复派发
  • _process_schedules:89)先把 next_run_at 推进到下一次(calculate_next_run_at), 再用 Celery group 批量投递 run_schedule_trigger,最后才 commit。
  • 有熔断:单次 tick 派发量超过 WORKFLOW_SCHEDULE_MAX_DISPATCH_PER_TICK 就 break(:45-50)。

执行侧 run_schedule_triggerapi/tasks/workflow_schedule_tasks.py:24,队列 schedule_executor) 做三件事:查计划、找租户 owner 当执行者(ScheduleService.get_tenant_ownertrigger/schedule_service.py:122)、以 ScheduleTriggerData(inputs={})trigger_workflow_async。注意 inputs 是空的——定时触发没有外部数据。

订阅刷新api/schedule/trigger_provider_refresh_task.py:49)是另一种批处理: _build_due_filter:22)把"凭据快过期"和"订阅快过期"两个条件 OR 起来(-1 表示永不过期, 被显式排除),分页扫描后用 _acquire_locks:36)在一次 Redis pipelineSET key 1 EX ttl NX 抢一批锁,只给抢到锁的那些投递 trigger_subscription_refresh。 锁键由 api/core/trigger/utils/locks.py:5build_trigger_refresh_lock_key 生成 (trigger_provider_refresh_lock:<tenant>_<sub>),下游任务在 api/tasks/trigger_subscription_refresh_tasks.py:86 先验锁存在、finally 里删锁。

4.4 汇流点与留痕

三条路最后都调 AsyncWorkflowService.trigger_workflow_asyncapi/services/async_workflow_service.py:56)。它做的事:校验 app/workflow → 按租户订阅等级挑队列 → 先写一条 WorkflowTriggerLog(status=PENDING):111-133)→ 投递 Celery,然后立刻返回。 入参统一是 TriggerData 家族(api/services/workflow/entities.py:28), 三个子类 WebhookTriggerData:44 / ScheduleTriggerData:51 / PluginTriggerData:71 各自钉死了 trigger_typetrigger_from

真正跑图在 api/tasks/async_workflow_tasks.py,注意 generate(...) 的两个参数 (:171-186):

generator.generate(
..., # 省略 app_model / workflow / user / args 等参数
root_node_id=trigger_data.root_node_id,
graph_engine_layers=[
TriggerPostLayer(cfs_plan_scheduler_entity, start_time, trigger_log.id),
],
)
  • root_node_id = 触发节点的 id。这就是"从哪个节点开跑"的落点,与 §3.1 的起点集合呼应。
  • TriggerPostLayer 是一个 GraphEngineLayer(机制见 04-execution-layers-and-tracing.md)。 它在 api/core/app/layers/trigger_post_layer.py:24,靠一张 _STATUS_MAP:29) 把四种终态事件映射成触发日志状态:
引擎事件触发日志状态
GraphRunSucceededEventSUCCEEDED
GraphRunFailedEventFAILED
GraphRunAbortedEventFAILED(并写 error
GraphRunPausedEventPAUSED

on_event:52)在命中这四种事件时开一个 session,用 SQLAlchemyWorkflowTriggerLogRepositoryapi/repositories/sqlalchemy_workflow_trigger_log_repository.py:18) 把 workflow_run_id、outputs、total_tokensfinished_at 回填。

PAUSED 这一档是给 human-in-the-loop 用的——所以 elapsed_time 用的是累加而不是覆盖 (trigger_post_layer.py:86-89):一次被人工审批打断再恢复的运行,耗时是几段之和。 暂停/恢复本身见 05-human-in-the-loop.md


5. 订阅生命周期:Dify 和插件之间的合约

5.1 六个动作

插件触发器的一切都围绕订阅(subscription):Dify 分配一个 endpoint URL,插件拿这个 URL 去第三方那儿注册 webhook;之后第三方的事件就会打到这个 URL 上。

TriggerManagerapi/core/trigger/trigger_manager.py:30)是这套合约的门面,全是 classmethod:

方法行号干什么
list_plugin_trigger_providers:49列出租户装了哪些触发器插件,包成 controller
get_trigger_provider:77取单个 provider,带请求级缓存
invoke_trigger_event:150把一次外部请求翻译成节点变量
subscribe_trigger:198让插件去第三方注册 webhook
unsubscribe_trigger:232反向注销
refresh_trigger:263续期即将过期的订阅

get_trigger_provider 的缓存写法值得抄:用 contexts.plugin_trigger_providers(ContextVar) 存字典,配一把 Lock先查 → 拿锁 → 再查一次(double check)→ 才去请求插件守护进程:92-118)。一次 HTTP 请求内多次要同一个 provider,只会真的问插件一次。

这六个方法自己不发 HTTP,都转给 PluginTriggerProviderControllerapi/core/trigger/provider.py:38):dispatch:268invoke_trigger_event:298subscribe_trigger:337unsubscribe_trigger:370refresh_trigger:396—— controller 再转给 PluginTriggerClient(§6)。

5.2 数据模型:声明侧和实例侧

api/core/trigger/entities/entities.py 里有两组模型,分清楚很重要:

声明侧(插件的 manifest 说"我能干什么"):

模型行号是什么
EventParameter:36一个事件参数的 schema(名字、类型、是否必填)
EventEntity:92一个事件(identity + parameters + output_schema)
SubscriptionConstructor:112建订阅需要用户填什么、要什么凭据、支不支持 OAuth
TriggerProviderEntity:138一个 provider 的全量声明

实例侧(这个租户真的建了一个订阅):

模型行号是什么
Subscription:155已建立的订阅:endpoint(Dify 分配的 URL)、parameterspropertiesexpires_at
SubscriptionBuilder:201半成品订阅,建的过程中暂存在 Redis 里
UnsubscribeResult:175退订结果

落库的那面在 api/models/trigger.pyTriggerSubscription:68WorkflowPluginTrigger:392 (哪个 app 的哪个节点订了哪个事件)、WorkflowWebhookTrigger:335AppTrigger:438(统一的启停状态)、 WorkflowSchedulePlan:487WorkflowTriggerLog:222

5.3 建一个订阅:为什么要个 "builder"

难点在于先有鸡还是先有蛋:第三方注册 webhook 时常常要立刻回调验证一次这个 URL, 但那时订阅还没建好。Dify 的解法是先发一个"半成品"占位。

用户在控制台填参数

├─ create_trigger_subscription_builder → 先分配 endpoint_id,builder 存 Redis(30min)
│ endpoint 立刻可用于验证回调
│ ↑ 验证请求走 process_builder_validation_endpoint

├─ (OAuth 时) 走 OAuthHandler 拿 credentials 回填进 builder

└─ build_trigger_subscription_builder
├─ credential_type == UNAUTHORIZED → 直接 add_trigger_subscription(手动模式)
└─ 否则 → TriggerManager.subscribe_trigger(endpoint=<URL>) 让插件去第三方注册
再 add_trigger_subscription,最后删掉 Redis 里的 builder

代码在 api/services/trigger/trigger_subscription_builder_service.py: 类 :35build_trigger_subscription_builder:107process_builder_validation_endpoint:445。 并发保护用 Redis 分布式锁 acquire_builder_lock:66,30 秒超时)。

/triggers/plugin/<endpoint_id> 这一个路由同时服务"正式订阅"和"builder 验证"两种流量—— controllers/trigger/trigger.py:25-28 用一条处理链依次试:

handling_chain = [
TriggerService.process_endpoint,
TriggerSubscriptionBuilderService.process_builder_validation_endpoint,
]

谁返回非 None 就用谁,都不认就 404。endpoint URL 由 api/core/trigger/utils/endpoint.py:11generate_plugin_trigger_endpoint_url 拼 (webhook 侧对应 :19generate_webhook_trigger_endpoint),基址来自配置 TRIGGER_URLapi/configs/feature/__init__.py:411)。

5.4 凭据加密与缓存

订阅里存着第三方的 API key 或 OAuth token,落库前要按 provider 声明的 schema 加密。 入口是 create_trigger_provider_encrypter_for_subscriptionapi/core/trigger/utils/encryption.py:54)——它把 TriggerProviderCredentialsCache:13) 和通用的 create_provider_encrypter 组装起来,缓存键按 tenant + provider + subscription_id 三元组隔离。

用的地方就在 §4.2 的同步段:trigger_service.py:101-111 先建 encrypter 再 encrypter.decrypt(subscription.credentials),解密后的明文只在这一次调用里传给插件。 masked_credentialsencryption.py:134)负责回给前端时打码。

CredentialType 三档(api/core/plugin/entities/plugin_daemon.py:251): API_KEY / OAUTH2 / UNAUTHORIZED。只有 API_KEY 可编辑、可校验 (is_editable:241 / is_validate_allowed:244);of() 做了 api_key/api-keyoauth/oauth2 的别名归一。

5.5 发布时的自动同步

用户在画布上加一个触发节点,数据库里的触发记录是发布时才对齐的。信号处理器 api/events/event_handlers/update_app_triggers_when_app_published_workflow_updated.py:15 监听 app_published_workflow_was_updated,用 TRIGGER_NODE_TYPES 从已发布图里捞出所有触发节点 (get_trigger_infos_from_workflow:99-110),和 AppTrigger 表做集合差分: 多出来的建、少掉的删。webhook 与插件触发还各有一套关系同步 (WebhookService.sync_webhook_relationshipswebhook_service.py:890TriggerService.sync_plugin_trigger_relationshipstrigger_service.py:153), 两者都用 Redis 缓存减少数据库读,且都限制每张图最多 5 个同类触发节点。


6. 节点能力从哪来:插件守护进程

6.1 一句话:Dify 主体不含工具和模型的实现

api/core/plugin/impl/ 下的每个文件都是一个 HTTP 客户端,共同基类 BasePluginClientapi/core/plugin/impl/base.py:96)。它们打的是同一个服务—— plugin daemon,地址来自 PLUGIN_DAEMON_URLbase.py:43), 用一个池化的 httpx.Client:62,keepalive 50 / 最大连接 100)。

八个客户端,各管一摊:

客户端位置管什么
PluginToolManagerimpl/tool.py:17工具的发现与调用
PluginModelClientimpl/model.py:30LLM/embedding/rerank/TTS/STT/moderation 的调用
PluginModelRuntimeimpl/model_runtime.py:106上者的门面,实现 ModelRuntime 协议
PluginDatasourceManagerimpl/datasource.py:23网页抓取、在线文档、网盘
PluginTriggerClientimpl/trigger.py:20§5 那套订阅动作
PluginAgentClientimpl/agent.py:14Agent 策略(ReAct/Function Calling 等)
PluginEndpointClientimpl/endpoint.py:8插件自己对外暴露的 HTTP 端点
OAuthHandlerimpl/oauth.py:11三方授权链路

6.2 请求形态高度统一

所有调用都长一个样:plugin/{tenant_id}/dispatch/<domain>/<action>, body 里嵌一层 data,插件身份放在 X-Plugin-ID 头。对比一下三个域:

干什么path出处
调工具plugin/{tenant_id}/dispatch/tool/invokeimpl/tool.py:85 invoke
调 LLMplugin/{tenant_id}/dispatch/llm/invokeimpl/model.py:168 invoke_llm
问"这请求算什么事件"plugin/{tenant_id}/dispatch/trigger/dispatch_eventimpl/trigger.py:157 dispatch_event
把请求翻成变量plugin/{tenant_id}/dispatch/trigger/invoke_eventimpl/trigger.py:81 invoke_trigger_event
拿 OAuth 授权地址plugin/{tenant_id}/dispatch/oauth/get_authorization_urlimpl/oauth.py:12

管理类(非 dispatch)走另一组:plugin/{tenant_id}/management/triggersimpl/trigger.py:37)、plugin/{tenant_id}/endpoint/setupimpl/endpoint.py:24)。

原始 HTTP 请求怎么塞进 JSON? 十六进制。触发器要把第三方的原始请求原样交给插件(插件可能要验签名), 所以 impl/trigger.py:110 写的是:

"raw_http_request": binascii.hexlify(serialize_request(request)).decode(),

6.3 流式协议:每行一个信封

工具和 LLM 都要流式返回。基类提供 _request_with_plugin_daemon_response_streambase.py:299):逐行读,每行用 PluginDaemonBasicResponse[T]api/core/plugin/entities/plugin_daemon.py:23)解析——统一信封{code, message, data}

  • code != 0 → 报错;特别地 code == -500 会把 message 再解一层 PluginDaemonError,交给 _handle_plugin_daemon_errorbase.py:341)翻译成 Dify 自己的异常类型——一长串 match 分支覆盖了限流、鉴权失败、CredentialsValidateFailedErrorTriggerInvokeErrorEventIgnoreError 等。
  • data is None → 直接报"空数据"。
  • 否则 yield rep.data

一次性返回的调用也复用同一条流,只是取第一帧就 return,取不到就抛 "No response received from plugin daemon"(impl/trigger.py:121-124 是典型写法)。

工具调用还在流外面套了一层 merge_blob_chunksimpl/tool.py:127), 把插件分片传回的二进制按 PLUGIN_MAX_FILE_SIZE 合并成完整文件。

6.4 反向:插件回头调 Dify

插件不只是被调用方。一个"写周报"的工具插件可能需要用租户配好的模型总结一段文字—— 这时它反过来打 Dify 的 inner API。

Dify api 进程 plugin daemon
──────────────── ────────────────

core/plugin/impl/* ──────正向:调能力────────> 工具/模型/触发器实现
(HTTP client) │
│ 反向:我也要用 Dify 的能力
controllers/inner_api/plugin/plugin.py <─────────────┘
@plugin_inner_api_only (校验 X-Inner-Api-Key)

└─> core/plugin/backwards_invocation/*
├─ model.py 用租户配置的模型
├─ tool.py 用另一个工具
├─ node.py 借用参数提取 / 问题分类节点
├─ app.py 调起一个完整的 Dify 应用
└─ encrypt.py 借 Dify 的加密能力存自己的配置
  • 入口全在 api/controllers/inner_api/plugin/plugin.py/invoke/llm:42)、 /invoke/tool:225)、/invoke/parameter-extractor:257)、 /invoke/question-classifier:290)、/invoke/app:323)、/invoke/encrypt:353)……
  • 鉴权是共享密钥:plugin_inner_api_onlyapi/controllers/inner_api/wraps.py:84) 比对 X-Inner-Api-Key 头与 INNER_API_KEY_FOR_PLUGIN,不匹配一律 404(不是 401——不暴露端点存在)。
  • 实现侧:BaseBackwardsInvocationbackwards_invocation/base.py:6)提供 convert_to_event_stream 把结果转成 SSE; PluginToolBackwardsInvocation.invoke_toolbackwards_invocation/tool.py:20)走的是 和普通工具节点同一条 ToolEngine.generic_invoke,并显式传 workflow_call_depth=1 防止无限套娃;PluginNodeBackwardsInvocationnode.py:15) 把两个 LLM 派生节点当纯函数暴露出去;PluginAppBackwardsInvocationapp.py:31)能拉起整个 app。

6.5 这些能力怎么进到节点里

第 03 章讲了 DifyNodeFactory 怎么构造节点对象。这里只补注入的是什么—— 它注入的正是通往插件守护进程的那条路。

工具节点:工厂给它一个 runtime(api/core/workflow/node_factory.py:488-491):

BuiltinNodeTypes.TOOL: lambda: {
"tool_file_manager": self._bound_tool_file_manager_factory(),
"runtime": self._tool_runtime,
},

self._tool_runtimeDifyToolNodeRuntimenode_factory.py:364 构造, 类在 api/core/workflow/node_runtime.py:524)。节点跑到时调它的 get_runtime:472), 里面是 ToolManager.get_workflow_tool_runtime(...);对插件类工具,最终拿到的是 PluginToolapi/core/tools/plugin_tool/tool.py:15),它的 _invoke:28)就是 PluginToolManager().invoke(...)——回到 §6.2 那个 dispatch/tool/invoke

LLM 节点:工厂在构造时就用 build_dify_model_accessapi/core/app/llm/model_access.py:117)拿到一对"凭据提供者 + 模型工厂" (node_factory.py:380);节点级则由 _build_model_instance_for_llm_nodenode_factory.py:706) 经 fetch_model_configmodel_access.py:169)造出 ModelInstanceapi/core/model_manager.py:36)。 build_dify_model_access 内部调的是 create_plugin_provider_managerapi/core/plugin/impl/model_runtime_factory.py:130,同文件 PluginModelAssembly:48 是装配枢纽), 所以 ModelInstance.model_type_instance 底下最终是 PluginModelRuntimePluginModelClient

链条完整地串起来是这样:

DifyNodeFactory ──注入──> ToolNodeRuntime / ModelInstance

v
ToolManager / ModelProviderFactory

v
PluginTool / PluginModelRuntime

v
PluginToolManager / PluginModelClient (core/plugin/impl/)
│ HTTP: plugin/{tenant}/dispatch/...
v
plugin daemon 进程

ModelInstance 还多做一件事:_round_robin_invokecore/model_manager.py:379)—— 一个 provider 配了多把 key 时轮询,某把出 InvokeRateLimitError 就换下一把。 这层容错在 Dify 侧,不在插件侧。


7. 巧妙之处(可以抄的)

  • 加起点只加一个字符串。 TRIGGER_NODE_TYPES* 展开进 _START_NODE_TYPESnode_factory.py:81-83),新增触发器类型时工厂、草稿变量、发布同步三处自动跟上。
  • 触发节点做成"只读变量池"。 三个 _run 都不含 IO,接收逻辑全在图外—— 于是同一个节点在正式运行、单节点调试、草稿预览下行为完全一致,因为它只认池子里有什么。
  • 同步段只认领、异步段才翻译。 process_endpointtrigger_service.py:77)在插件说完 "有哪些事件"后立刻回响应,把 N 次"翻译成变量"的插件调用甩进队列——第三方不会因为你订阅了 10 张图而超时。
  • 跨进程传原始请求用对象存储 + 短 id。 TriggerHttpRequestCachingServicetrigger_request_service.py:11)把 Flask Request 序列化落盘,Celery 消息里只带 request_id, 队列不被大 payload 撑爆。
  • 轮询用 skip_locked 而不是分布式锁。 _fetch_due_schedulesworkflow_schedule_task.py:82)靠数据库行锁天然去重,多 poller 并存零协调成本。
  • 批量抢锁走一次 pipeline。 _acquire_lockstrigger_provider_refresh_task.py:36) 把一批 SET NX EX 压进一个 pipeline,一次往返决定这一页里哪些该刷新。
  • 配额预留—提交—退还三态。 触发路径上每一处扣费都配了 commit() / refund()webhook_service.py:828-840trigger_processing_tasks.py:304-400),插件说"忽略这次事件" 同样退款。
  • 限流是"标记全租户"而不是"每次都算"。 mark_tenant_triggers_rate_limitedapp_trigger_service.py:24)一次 UPDATE 把租户所有触发器改成 RATE_LIMITED, 后续请求在查询阶段就被挡。
  • 插件失败也要留下一条可见的运行。 _record_trigger_failure_logtrigger_processing_tasks.py:126)——最难查的 bug 是"事件好像没到",这条日志把它变成"到了但失败了"。
  • 反向调用鉴权失败返 404 不返 401。 wraps.py:84-93,连端点是否存在都不告诉你。

8. 边界与局限(诚实)

  • 一张图只支持一个定时计划。 extract_schedule_configschedule_service.py:198)在注释里 写明"只返回找到的第一个 schedule 节点,不支持多个"(:227-229)。
  • 每张图最多 5 个同类触发节点。 webhook 和插件触发的同步方法 docstring 里都硬写着 Maximum 5 ... per workflowwebhook_service.py:903trigger_service.py:166), 常量是 MAX_PLUGIN_TRIGGER_NODES_PER_WORKFLOWtrigger_service.py:41)。
  • webhook 只有异步模式。 WebhookData.SyncMode 只定义了一个成员,而且写成 SYNC = "async" # only supporttrigger_webhook/entities.py:90-91)——命名和取值矛盾, 实际含义是"只支持异步"。节点默认配置里的 "async_mode": Truetimeout: 30trigger_webhook/node.py:45,48)目前也没有对应的同步等待实现。
  • refresh_trigger 的返回值没落库。 trigger_manager.py:282 留着 # TODO you should update the subscription using the return value of the refresh_trigger—— 刷新后插件返回的新 Subscription 当前被丢弃。
  • 插件触发节点的事件参数不能引用变量。 见 §3.3,resolve_parameterstrigger_plugin/entities.py:76-78)遇到非 constant 直接抛错。
  • webhook 调试端点会静默丢事件。 没有活跃监听器时 handle_webhook_debugcontrollers/trigger/webhook.py:133-144)只记 warning 并返回错误,请求本身就没了。
  • 不覆盖的内容: 节点对象怎么被构造(见 03)、 Layer 机制本身(见 04)、 暂停恢复(见 05)。

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

触发节点与起点判定

主题文件符号
三个触发类型常量api/core/trigger/constants.pyTRIGGER_NODE_TYPESis_trigger_node_type
起点集合api/core/workflow/node_factory.py_START_NODE_TYPESis_start_node_type
webhook 节点api/core/workflow/nodes/trigger_webhook/node.pyTriggerWebhookNode_extract_configured_outputs
webhook 节点配置api/core/workflow/nodes/trigger_webhook/entities.pyWebhookDataMethodContentType
定时节点api/core/workflow/nodes/trigger_schedule/trigger_schedule_node.pyTriggerScheduleNode
定时节点配置api/core/workflow/nodes/trigger_schedule/entities.pyTriggerScheduleNodeDataVisualConfigScheduleConfig
插件事件节点api/core/workflow/nodes/trigger_plugin/trigger_event_node.pyTriggerEventNode
插件事件节点配置api/core/workflow/nodes/trigger_plugin/entities.pyTriggerEventNodeDataresolve_parameters

订阅与外部事件

主题文件符号
订阅门面api/core/trigger/trigger_manager.pyTriggerManagersubscribe_trigger / invoke_trigger_event / refresh_trigger
provider 控制器api/core/trigger/provider.pyPluginTriggerProviderControllerdispatch
声明与实例模型api/core/trigger/entities/entities.pyEventEntitySubscriptionSubscriptionBuilder
endpoint URL 拼装api/core/trigger/utils/endpoint.pygenerate_plugin_trigger_endpoint_urlgenerate_webhook_trigger_endpoint
刷新锁api/core/trigger/utils/locks.pybuild_trigger_refresh_lock_key
凭据加解密api/core/trigger/utils/encryption.pycreate_trigger_provider_encrypter_for_subscriptionmasked_credentials
webhook 全流程api/services/trigger/webhook_service.pyWebhookServicetrigger_workflow_executionsync_webhook_relationships
插件事件同步段api/services/trigger/trigger_service.pyTriggerService.process_endpoint
原始请求存取api/services/trigger/trigger_request_service.pyTriggerHttpRequestCachingService
订阅建造api/services/trigger/trigger_subscription_builder_service.pybuild_trigger_subscription_builderprocess_builder_validation_endpoint
订阅方查询api/services/trigger/trigger_subscription_operator_service.pyget_subscriber_triggers
限流标记api/services/trigger/app_trigger_service.pyAppTriggerService.mark_tenant_triggers_rate_limited
排程配置api/services/trigger/schedule_service.pyScheduleService.visual_to_cronextract_schedule_config
事件路由api/controllers/trigger/trigger.py / webhook.pytrigger_endpointhandle_webhookhandle_webhook_debug
异步派发api/tasks/trigger_processing_tasks.pydispatch_triggered_workflowdispatch_triggered_workflows_async
beat:排程轮询api/schedule/workflow_schedule_task.pypoll_workflow_schedules_fetch_due_schedules
beat:订阅刷新api/schedule/trigger_provider_refresh_task.pytrigger_provider_refresh_acquire_locks
beat 注册api/extensions/ext_celery.pybeat_schedule["workflow_schedule_task"]beat_schedule["trigger_provider_refresh"]
发布时同步api/events/event_handlers/update_app_triggers_when_app_published_workflow_updated.pyget_trigger_infos_from_workflow
运行后回填api/core/app/layers/trigger_post_layer.pyTriggerPostLayer_STATUS_MAP
触发日志仓储api/repositories/sqlalchemy_workflow_trigger_log_repository.pySQLAlchemyWorkflowTriggerLogRepository
触发相关表api/models/trigger.pyTriggerSubscriptionWorkflowWebhookTriggerAppTriggerWorkflowSchedulePlanWorkflowTriggerLog

插件运行时

主题文件符号
HTTP 基类与流协议api/core/plugin/impl/base.pyBasePluginClient_request_with_plugin_daemon_response_stream
工具调用api/core/plugin/impl/tool.pyPluginToolManager.invoke
模型调用api/core/plugin/impl/model.pyPluginModelClient.invoke_llm
模型运行时门面api/core/plugin/impl/model_runtime.pyPluginModelRuntime
运行时装配api/core/plugin/impl/model_runtime_factory.pyPluginModelAssemblycreate_plugin_model_runtime
数据源api/core/plugin/impl/datasource.pyPluginDatasourceManager
触发器客户端api/core/plugin/impl/trigger.pyPluginTriggerClient.dispatch_eventinvoke_trigger_event
Agent 策略api/core/plugin/impl/agent.pyPluginAgentClient.invoke
插件端点api/core/plugin/impl/endpoint.pyPluginEndpointClient.create_endpoint
OAuthapi/core/plugin/impl/oauth.pyOAuthHandler.get_authorization_urlrefresh_credentials
守护进程信封与类型api/core/plugin/entities/plugin_daemon.pyPluginDaemonBasicResponseCredentialTypePluginTriggerProviderEntity
插件安装与缓存api/core/plugin/plugin_service.pyPluginService.fetch_plugin_model_providers
反向调用入口api/controllers/inner_api/plugin/plugin.pyPluginInvokeLLMApiPluginInvokeToolApi
反向调用鉴权api/controllers/inner_api/wraps.pyplugin_inner_api_only
反向调用实现api/core/plugin/backwards_invocation/BaseBackwardsInvocationPluginToolBackwardsInvocationPluginNodeBackwardsInvocationPluginAppBackwardsInvocation
模型实例与轮询api/core/model_manager.pyModelInstance_round_robin_invoke
插件工具适配api/core/tools/plugin_tool/tool.pyPluginTool._invoke
工具注入点api/core/workflow/node_runtime.pyDifyToolNodeRuntime.get_runtime