跳到主要内容

数据截至 (上游 commit b78a3462c9a6)

一次运行的生命周期:从 HTTP 请求到 SSE 流

30 秒导读: 你在 Dify 画布上点一次「运行」,浏览器发出的是一个普通 HTTP POST,收回来的却是一条持续几十秒的 SSE 事件流。这中间 Dify 做了一件核心的事:把「执行」和「回话」拆到两个执行体上——执行体在后台跑图、往队列里丢事件,HTTP 这一侧只负责从队列里捞事件、翻译成 data: {...} 往下吐。本章只讲这条骨架,不进节点内部。


1. 这章讲什么

一句话:用户点一次运行,后端从收到请求到吐完最后一个 SSE chunk,中间经过了哪些部件、按什么顺序。

读完你应该能回答四个问题:

问题本章第几节
AppMode 的八个取值怎么分流?谁挡住并发?§3
谁真正在跑工作流?为什么不能在 HTTP 线程里直接跑?§4、§5、§6
「停止」按钮按下去,凭什么能把跑到一半的执行掐掉?§7
引擎里的一个事件,怎么变成浏览器收到的一行 data:§8、§9

本章不覆盖(在兄弟章里):


2. 顶层全景:一次运行分成五段路

先给一张最粗的图。从上往下读,就是一次请求的时间顺序;左边一列是「回话」的路,右边一列是「执行」的路,中间那条竖线是两条路唯一的接触面——队列。

浏览器:POST .../workflows/draft/run


① 入口分发 AppGenerateService.generate ← 挑模式 + 限流 + 计费


② 起摊子 XxxAppGenerator.generate/_generate

┌──────────────────┴──────────────────┐
│ │
(回话侧 / 调用方线程) (执行侧 / 工作线程或 Celery worker)
│ │
▼ ▼
③ 消费 GenerateTaskPipeline ③' 生产 XxxAppRunner.run()
.process() → WorkflowEntry → GraphEngine
│ │
│ ◄──────── AppQueueManager ─────────┤ 引擎事件 publish 进队列
│ (queue.Queue) │
▼ ▼
④ 翻译 Queue*Event → *StreamResponse 跑完发终止事件 → stop_listen()


⑤ 出口 ResponseConverter → "data: {...}\n\n" → Flask SSE Response

五段路各归谁管:

干什么主文件
① 入口分发按 app 模式选 Generator,做并发限流与配额api/services/app_generate_service.py
② 起摊子组装运行实体、仓储、队列管理器,起执行体api/core/app/apps/*/app_generator.py
③' 生产建图、跑 GraphEngine、把引擎事件塞进队列api/core/app/apps/*/app_runner.py
③ 消费监听队列,逐事件产出流式响应对象api/core/app/apps/*/generate_task_pipeline.py
⑤ 出口响应对象 → dict → SSE 文本行api/core/app/apps/*/generate_response_converter.py

3. 第一关:模式分发与并发限流

3.1 一个 match-case 管住所有入口

所有 HTTP 控制器(console / web / service_api / openapi)最后都汇到同一个静态方法:AppGenerateService.generateapi/services/app_generate_service.py:88,带 @trace_span(AppGenerateHandler) 装饰)。

分发之前先做一次模式修正:如果这个 App 挂了旧版 agent 标记,就强行按 agent-chat 走(api/services/app_generate_service.py:170-174effective_mode)。修正后进入 match effective_mode:179)。

AppMode 一共八个取值(api/models/model.py:373AppMode):completion / workflow / chat / advanced-chat / agent-chat / agent / channel / rag-pipeline。这个 match 显式列出其中六个,另外两个(channelrag-pipeline)落进 case _

模式走哪个 Generator执行位置源码位置
completionCompletionAppGenerator同进程工作线程app_generate_service.py:176
chatChatAppGenerator同进程工作线程app_generate_service.py:218
agent-chatAgentChatAppGenerator同进程工作线程app_generate_service.py:190
agentAgentAppGenerator同进程工作线程app_generate_service.py:204
advanced-chatAdvancedChatAppGenerator流式走 Celery;阻塞走同进程app_generate_service.py:232-289
workflowWorkflowAppGenerator流式走 Celery;阻塞走同进程app_generate_service.py:290-342
channelcase _,报错app_generate_service.py:343-344
rag-pipelinecase _,另有独立入口同进程工作线程services/rag_pipeline/pipeline_generate_service.py:19

三条值得记住的边界:

  • 落进 case _ 的两个模式命运不同。 channelrag-pipeline 在这里都会 raise ValueErrorapp_generate_service.py:343-344),但 rag-pipeline 另有入口 PipelineGenerateService.generate,且复用同一套四件套(见 §10);channel 在整个 AppGenerateService 里都被归到报错侧——单迭代/单循环那两个 match 甚至给它写了显式的拒绝分支(:411-412:456-457)。
  • agentagent-chat 是两条独立分支。 agent-chat 是旧的 ReAct 应用,agent 是 Dify Agent 运行时支撑的新应用类型(枚举定义处的注释写明二者 distinct,api/models/model.py:380-383)。但注意 §上面那次模式修正:挂了旧版 agent 标记的 App 会被强行拨回 agent-chat,所以 case AppMode.AGENT 只接得到真正的新型 Agent 应用。
  • advanced-chat / workflow 的流式和阻塞是两条不同的物理路径,不只是返回类型不同。这是这份代码里最容易看漏的一点,§4 单独讲。

3.2 并发闸门:_get_max_active_requests

限流的额度算法只有四行,规则是「取两个上限里更小的那个,0 代表不限」(api/services/app_generate_service.py:360-365_get_max_active_requests):

app_limit = app.max_active_requests or dify_config.APP_DEFAULT_ACTIVE_REQUESTS
config_limit = dify_config.APP_MAX_ACTIVE_REQUESTS
limits = [limit for limit in [app_limit, config_limit] if limit > 0]
return min(limits) if limits else 0

这段说的是:App 自己的配额优先,缺省回落到全局默认;两者都是 0(默认值,见 api/configs/feature/__init__.py:89-96)时完全不限流。

拿到额度后交给 RateLimit

动作做了什么位置
rate_limit.enter(request_id)读 Redis hash 里在途请求数,超了抛 AppInvokeQuotaExceededErrorcore/app/features/rate_limiting/rate_limit.py:73
rate_limit.exit(request_id)从 hash 里删掉这个请求rate_limit.py:90
rate_limit.generate(gen, request_id)把生成器包成 RateLimitGeneratorrate_limit.py:102

关键设计:流式响应的「归还名额」时机不是函数返回,而是生成器关闭。 阻塞模式在 finally: if not streaming: rate_limit.exit(...) 里归还(app_generate_service.py:152-154);流式模式则由 RateLimitGenerator.close() 在迭代结束或异常时归还(rate_limit.py:140-145)。名额的生命周期跟着 SSE 连接走,而不是跟着 HTTP handler 走。

计费配额同理:QuotaService.reserve 先占,成功后 commit,异常路径 refund——限流与配额如今统一收进 _run_with_guardrails 这个包装器(app_generate_service.py:125-154,云上部署才真正占用配额)。


4. 两个执行位置:同进程线程 vs Celery worker

先说结论:「执行」不一定在 web 进程里。

路径谁跑 Runner事件怎么回到 HTTP 侧
completion / chat / agent-chat / agent(全部)
advanced-chat / workflow(阻塞
同进程的 threading.Thread进程内 queue.Queue
advanced-chat / workflow(流式Celery worker 进程Redis 广播 topic(跨进程)

流式分支的代码长这样(api/services/app_generate_service.py:236-265 为 advanced-chat,:297-329 为 workflow):把参数序列化成 AppExecutionParamsapi/tasks/app_generate/workflow_execute_task.py:85,构造器 new:98),HTTP 侧只订阅事件流:

on_subscribe = cls._build_streaming_task_on_subscribe(on_subscribe)
return rate_limit.generate(
WorkflowAppGenerator.convert_to_event_stream(
MessageBasedAppGenerator.retrieve_events(
AppMode.WORKFLOW, payload.workflow_run_id, on_subscribe=on_subscribe,
),
),
request_id,
)

这段的意思:HTTP 线程不跑工作流,它只是订阅 channel:{app_mode}:{workflow_run_id} 这个 topic(key 生成在 api/core/app/apps/message_based_app_generator.py:308_make_channel_key),谁往里发事件就往下吐谁。

4.1 「什么时候真正开跑」被订阅时机门控

_build_streaming_task_on_subscribeapi/services/app_generate_service.py:44-84)解决一个很实际的竞态:Celery 任务如果先跑完再有人订阅,pub/sub 模式下前面的事件就丢了。

PUBSUB_REDIS_CHANNEL_TYPE == "streams"
└─► 立刻 delay() 开跑(事件在 Redis Stream 里可回读,晚到的订阅者能补看)

PUBSUB_REDIS_CHANNEL_TYPE == "pubsub" / "sharded"(至多一次投递)
├─► 等第一个订阅者到达再 delay()
└─► 同时挂一个 200ms 兜底定时器:客户端始终不连也照跑

SSE_TASK_START_FALLBACK_MS = 200:34),通道类型来自 PUBSUB_REDIS_CHANNEL_TYPEapi/configs/middleware/cache/redis_pubsub_config.py:49,实际选型在 api/extensions/ext_redis.py:495get_pubsub_broadcast_channel)。

4.2 Celery 侧做的事

workflow_based_app_execution_taskapi/tasks/app_generate/workflow_execute_task.py:472)反序列化参数后交给 _AppRunner.run:152):

  1. _setup_flask_context:146)进 flask_app.app_context()set_login_user(user)
  2. _run_app:200)调同一个 AdvancedChatAppGenerator().generate(...) / WorkflowAppGenerator().generate(...)——也就是说,Celery worker 里跑的还是那套四件套,只是宿主换了进程;
  3. 拿到生成器后交给 _publish_streaming_response:342),逐事件 topic.publish(payload.encode()):426)。

_publish_streaming_response 有一处很值得学的兜底:它跟踪 workflow_started 与终止事件是否发出过,一旦生成器抛异常或没发终止事件就结束,它会补发一个 WorkflowFinishStreamResponse(status=FAILED):432-454)。目的是保证 SSE 消费者永远等得到终局事件,不会挂在半空。

4.3 HTTP 侧怎么读

retrieve_eventsapi/core/app/apps/message_based_app_generator.py:318)把订阅逻辑交给 stream_topic_eventsapi/core/app/apps/streaming_utils.py:13):

  • 一进来先 yield StreamEvent.PING.value:25),避免连接长时间处于 pending;
  • on_subscribe 在订阅真正生效后才回调(:33-34)——这正是 §4.1 门控的落点;
  • 默认终止事件是 workflow_finishedworkflow_paused:61-63),收到即 return
  • 长时间无消息(默认 idle_timeout=300message_based_app_generator.py:322)也退出。

5. 四件套协作模型

不管跑在哪个进程,_generate 内部永远是这四个部件在配合。

5.1 各自的职责

部件一句话职责代表文件关键符号
Generator组装参数、建仓储、起执行体、返回响应core/app/apps/workflow/app_generator.pyWorkflowAppGenerator._generate:315
Runner建图、跑 GraphEngine、把引擎事件转成队列事件core/app/apps/workflow/app_runner.pyrunapp_runner.py:74
QueueManager线程间信箱 + 停止标志core/app/apps/base_app_queue_manager.pylistenbase_app_queue_manager.py:64)/ publish:152
GenerateTaskPipeline消费队列,翻译成流式响应对象core/app/apps/workflow/generate_task_pipeline.pyprocessgenerate_task_pipeline.py:126

一句话记住它们的关系:Generator 是导演,Runner 是演员,QueueManager 是传声筒,Pipeline 是字幕组。

5.2 骨架长什么样(示意)

下面这段把 _generate 的核心骨架剥干净,帮你先建立直觉。

# 示意,非源码
def _generate(...):
queue_manager = WorkflowAppQueueManager(task_id=..., user_id=..., ...) # 传声筒

context = contextvars.copy_context() # 复制上下文,准备搬到新线程
db.session.close() # 先还掉数据库连接,后面要跑很久

threading.Thread( # 演员进后台开跑
target=self._generate_worker,
kwargs={"flask_app": current_app._get_current_object(),
"queue_manager": queue_manager, "context": context, ...},
).start()

response = self._handle_response(queue_manager=queue_manager, stream=streaming) # 字幕组开工
return WorkflowAppGenerateResponseConverter.convert(response, invoke_from)

重点看两处:起线程和 _handle_response紧挨着的——线程一启动,调用方立刻转去监听队列,两边同时跑。

5.3 真实源码对照

WorkflowAppGenerator._generateapi/core/app/apps/workflow/app_generator.py:321)的真实顺序:

步骤行号干什么
建队列管理器:352WorkflowAppQueueManager(...)
挂暂停持久化 Layer:359-366pause_state_config 才加(详见 05
复制 contextvars:369contextvars.copy_context()
释放 DB 连接:372db.session.close()
起线程:374-390threading.Thread(target=self._generate_worker, ...)
监听队列:395self._handle_response(...)
转换出口格式:404WorkflowAppGenerateResponseConverter.convert(...)

advanced-chat 版本结构一样(api/core/app/apps/advanced_chat/app_generator.py:508,起线程在 :553,监听在 :580),只多两件事:先建会话与消息记录(:517_init_generate_records),以及在关闭 session 前把 ORM 对象拍成快照:574-577):

workflow_snapshot = WorkflowSnapshot.from_workflow(workflow)
conversation_snapshot = ConversationSnapshot.from_conversation(conversation)
message_snapshot = MessageSnapshot.from_message(message)
db.session.close()

这三行是为了让下游 pipeline 只拿标量字段,不再碰已经关掉 session 的 ORM 实例。

Runner 那一侧的收尾同样简单——WorkflowAppRunner.run 建完 WorkflowEntryapi/core/app/apps/workflow/app_runner.py:171)后就是一个 for 循环(:186-189):

generator = workflow_entry.run()
for event in generator:
self._handle_event(workflow_entry, event)

_handle_eventapi/core/app/apps/workflow_app_runner.py:409)是一张巨大的 match,把 graphon 的 GraphEngineEvent 一一映射成 Queue*Event,最后统一经 _publish_event:712)投进队列。


6. 为什么必须开工作线程,以及上下文怎么搬过去

6.1 为什么不能在 HTTP 线程里直接跑

因为流式响应必须边跑边吐。如果在同一个线程里先把工作流跑完再返回,SSE 就退化成「等 60 秒然后一次性收到全部」。

所以 Dify 把它拆成经典的生产者-消费者:

HTTP 线程(消费者) 工作线程(生产者)
│ │
│ listen() 阻塞等 │ 跑节点…
│ ◄──── queue.put(事件) ─────┤
│ yield SSE chunk │ 跑节点…
│ ◄──── queue.put(事件) ─────┤
│ │ 跑完 → put(终止事件) → stop_listen()
│ ◄──── queue.put(None) ─────┘
▼ listen() 收到 None,break

6.2 新线程会丢三样东西

Python 的线程不会自动继承 Flask 的请求/应用上下文,也不继承 contextvars。丢掉的东西具体是:

丢的东西后果怎么补
Flask app contextcurrent_appdb 扩展不可用flask_app.app_context()
g._login_userflask-login 的 current_user 变空进新 context 后手动写回
contextvars插件工具提供者等上下文变量丢失逐个 var.set(val)

三件事被打包进一个上下文管理器 preserve_flask_contextsapi/libs/flask_utils.py:13),逻辑就三步:先把传进来的 contextvars 全部 set:46-48),保存 g._login_user:51-54),进 flask_app.app_context() 后再把 user 写回(:57-61)。

工作线程入口的第一行就是它(api/core/app/apps/workflow/app_generator.py:638):

with preserve_flask_contexts(flask_app, context_vars=context):

advanced-chat 版本在 api/core/app/apps/advanced_chat/app_generator.py:673,pipeline 版本在 api/core/app/apps/pipeline/pipeline_generator.py:319

chat / completion / agent-chat 用的是一个等价变体:@copy_current_request_context 装饰 + context.run(...)api/core/app/apps/chat/app_generator.py:219-230)——它们需要的是请求上下文,因为还要读请求级信息。

6.3 跨线程的硬约束:不许传 ORM 对象

工作线程里所有异常都被兜住并转成队列事件,而不是往上抛(app_generator.py:643-661):GenerateTaskStoppedError 只记日志,其余一律 queue_manager.publish_error(e, PublishFrom.APPLICATION_MANAGER)。原因很直白——调用方线程已经在 listen() 里等着了,抛出去没人接,只有变成事件才能传回去。

与此配套的还有一道运行时护栏:publish 每次都会先递归扫一遍事件负载,发现 SQLAlchemy 模型实例就直接 TypeErrorapi/core/app/apps/base_app_queue_manager.py:235-249_check_for_sqlalchemy_models):

"Critical Error: Passing SQLAlchemy Model instances that cause thread safety issues is not allowed."

这解释了 §5.3 里那三个 Snapshot 类为什么存在。


7. 队列与停止

7.1 listen():一个带三重超时的循环

AppQueueManager.listenapi/core/app/apps/base_app_queue_manager.py:64)是消费侧的心脏,一次循环干四件事:

动作行号说明
self._q.get(timeout=1):671 秒轮询,避免死等
message is Nonebreak:68-69None 是关闭哨兵
超时或被停 → 补发 QueueStopEvent:76-81总时长上限 APP_MAX_EXECUTION_TIME(默认 1200 秒,api/configs/feature/__init__.py:85
每 10 秒发一次 QueuePingEvent:83-85心跳,防连接被中间层掐断

生产侧 publish:152)先过 ORM 护栏,再交给子类 _publish。子类做的额外一件事是识别终止/暂停事件并主动关闭队列api/core/app/apps/workflow/app_queue_manager.py:35-48):

if isinstance(event, QueueWorkflowPausedEvent):
self.stop_listen(execution_state=AppExecutionState.PAUSED)
elif isinstance(event, QueueStopEvent | QueueErrorEvent | QueueMessageEndEvent
| QueueWorkflowSucceededEvent | QueueWorkflowFailedEvent
| QueueWorkflowPartialSuccessEvent):
self.stop_listen(execution_state=AppExecutionState.TERMINAL)

stop_listenbase_app_queue_manager.py:99-110)做四件事:按传入的 execution_state 把本次监听段标记成 PAUSED/TERMINAL、置位监听段完成标志、删掉 Redis 归属键,最后 self._q.put(None) 放哨兵——这就是 listen() 退出的正常路径。

消息类应用的 _publish 多一条规则(api/core/app/apps/message_based_app_queue_manager.py:59-62):来自 APPLICATION_MANAGER 的发布,如果发现任务已被停止,非 advanced-chat 模式直接 raise GenerateTaskStoppedError(),让工作线程当场退出。

7.2 停止:两个 Redis 键 + 两套机制

先看两个键(base_app_queue_manager.py:218-233):

什么时候写TTL用途
generate_task_belong:{task_id}队列管理器构造时(:511800s记录任务属于哪个用户,用于停止时的权限校验
generate_task_stopped:{task_id}调用 set_stop_flag 时(:191600s停止标志本身

set_stop_flag:177)先读归属键,比对 account-{id} / end-user-{id} 前缀,不匹配就静默返回;另有一个绕过校验的 set_stop_flag_no_user_check:194,内部转手给 set_app_task_stop_flag)。

读侧 _is_stopped:205)被 @cachedmethod 包着,走 1 秒 TTL 的 TTLCache:55)——listen() 每秒轮询一次,如果每次都打 Redis 会很费,缓存 1 秒刚好。

停止按钮打下来时,两套机制同时发api/controllers/console/app/workflow.py:1192-1203):

POST /apps/{id}/workflow-runs/tasks/{task_id}/stop

├─► AppQueueManager.set_stop_flag_no_user_check(task_id) (旧机制:Redis 标志位,
│ 下一轮 listen 轮询时生效)

└─► GraphEngineManager(redis_client).send_stop_command(task_id) (新机制:命令通道,
引擎主动中断当前执行)

源码注释直说了这是「为向后兼容,两条都发」。旧机制的粒度是「等下一次轮询」,新机制走的是 Runner 建的 Redis 命令通道(api/core/app/apps/workflow/app_runner.py:156-164,key 由 app_task_command_channel_key 生成,格式仍是 workflow:{task_id}:commands)。

本节只讲「按钮这一侧怎么发」。命令通道的另一半——Redis 与内存两种通道怎么选、除了用户之外还有哪些 Layer 会往通道里发 Abort/Pause——见 04 执行期横切 §4


8. 事件如何变成 SSE:四层翻译

一个引擎里的「节点跑完了」,要经过四次变形才成为浏览器上的一行文本。

数据形态谁做的位置
1GraphEngineEvent(graphon 的引擎事件)WorkflowBasedAppRunner._handle_eventcore/app/apps/workflow_app_runner.py:409
2Queue*Event(进队列的应用事件)_publish_eventAppQueueManager.publishworkflow_app_runner.py:728
3*StreamResponse(结构化响应对象)GenerateTaskPipeline._process_stream_responseworkflow/generate_task_pipeline.py:696
4"data: {...}\n\n"(SSE 文本行)BaseAppGenerator.convert_to_event_streamcore/app/apps/base_app_generator.py:313

8.1 第 3 层:pipeline 的双层分发

_process_stream_responseapi/core/app/apps/workflow/generate_task_pipeline.py:696)就是消费 listen() 的那个 for 循环:

for queue_message in self._base_task_pipeline.queue_manager.listen():
event = queue_message.event
match event:
case QueueWorkflowStartedEvent(): ...
case QueueErrorEvent():
yield from self._handle_error_event(event); break
case QueueWorkflowPausedEvent():
yield from self._handle_workflow_paused_event(event); break
case _:
...self._dispatch_event(event, ...)

设计上分成两层:

  • match 里显式列出的:708-731)都是break 循环的终局事件:错误、失败、暂停、停止。它们决定这一趟流什么时候结束。
  • 其余事件交给 _dispatch_event:645),它查一张 dict[事件类型, 处理函数] 表(_get_event_handlers:614-643),命中就调,不命中就静默丢弃(:693-694)。节点开始/成功/重试、迭代、循环、agent 日志、人工输入表单都在这张表里。

advanced-chat 的 pipeline 结构完全对称(api/core/app/apps/advanced_chat/generate_task_pipeline.py:993),差别是它把 QueueWorkflowSucceededEvent / QueueWorkflowPartialSuccessEvent 也列进了 break 分支(:1009-1015)——因为聊天模式在工作流成功后还要落消息记录、生成会话名。

单个 handler 通常很短,例如文本增量(workflow/generate_task_pipeline.py:557-574):拿到 event.text,顺手喂给 TTS 发布器,然后 yield 一个 TextChunkStreamResponse(构造在 :778)。

8.2 第 4 层:两种 SSE 行

convert_to_event_streamapi/core/app/apps/base_app_generator.py:313-328)只有一个判断:

if isinstance(message, Mapping | dict):
yield f"data: {orjson_dumps(message)}\n\n"
else:
yield f"event: {message}\n\n"

字典 → data: 行;裸字符串 → event: 行。后者的唯一来源是 ping:converter 遇到 PingStreamResponseyield "ping"api/core/app/apps/workflow/generate_response_converter.py:59-61)。


9. 两种返回形态:blocking 与 streaming

分叉点只有一个——GenerateTaskPipeline.process()api/core/app/apps/workflow/generate_task_pipeline.py:126-139):

process()

stream=True ───┴─── stream=False
│ │
▼ ▼
_to_stream_response _to_blocking_response
(:231) (:141)
│ │
逐个包成 遍历整条流,只截取终局:
WorkflowAppStreamResponse · ErrorStreamResponse → raise
往外 yield · WorkflowPauseStreamResponse → PausedBlockingResponse
· WorkflowFinishStreamResponse → BlockingResponse
· 其余 → continue(丢弃)

关键认知:阻塞模式并不是「另一套执行逻辑」,它跑的是同一条流,只是把中间事件全丢掉、只留最后一个(:193-194case _: continue)。如果流结束了一个终局事件都没等到,就 raise ValueError("queue listening stopped unexpectedly."):199)。

再往外一层是详略分叉,由 invoke_from 决定(api/core/app/apps/base_app_generate_response_converter.py:24-44):

调用来源走哪套差别
DEBUGGERSERVICE_APIconvert_*_full_response节点事件带完整 inputs/outputs/process_data
其余(WEB_APPEXPLORE 等)convert_*_simple_response节点事件走 to_ignore_detail_dict(),剥掉细节

full / simple 的实际差异就在一个 case 上(api/core/app/apps/workflow/generate_response_converter.py:103-108):simple 版对 NodeStartStreamResponse | NodeFinishStreamResponseto_ignore_detail_dict(),其余分支两版一模一样。画布调试看得到节点内部,公开 WebApp 看不到,就是这一行决定的。

最后一步是控制器统一收口:compact_generate_responseapi/libs/helper.py:412-431)——Mapping 就返回 application/json,否则返回 mimetype="text/event-stream" 的流式 Response


10. RAG_PIPELINE:换了入口,没换骨架

rag-pipeline 模式不走 AppGenerateService,入口是 PipelineGenerateService.generateapi/services/rag_pipeline/pipeline_generate_service.py:19)。但打开 PipelineGenerator._generateapi/core/app/apps/pipeline/pipeline_generator.py:289)会发现四件套原封不动:

部件复用情况位置
QueueManager自己的 PipelineQueueManager,基类同一个pipeline_generator.py:324
工作线程同样 threading.Thread + _generate_worker:325-339:550
Pipeline直接复用 WorkflowAppGenerateTaskPipeline:653
Converter直接复用 WorkflowAppGenerateResponseConverter:355

唯一结构性差异:PipelineGenerator._generate整个方法体包在 preserve_flask_contexts 里(:311),而 workflow / advanced-chat 只在工作线程入口用它。


11. 边界与容易踩的坑

  • _generate_worker 里的异常不会传播给调用方。 它们被转成 QueueErrorEventworkflow/app_generator.py:692-705)。调试时不要只看 HTTP 500,要看 SSE 流里的 error 事件。
  • 流式的 advanced-chat / workflow 需要 Celery worker 在线。 web 进程只订阅不执行;worker 没起来,客户端会一直等到 idle_timeout(300 秒)后静默断开。
  • 停止不是即时的。 旧机制最坏要等一轮 listen() 轮询(1 秒);新命令通道由引擎侧决定何时响应。
  • generate_task_stopped 键 TTL 只有 600 秒base_app_queue_manager.py:191),比 APP_MAX_EXECUTION_TIME(1200 秒)短。对超长执行,标志位可能在执行结束前就过期。
  • 队列事件里绝不能夹带 ORM 对象,否则运行时直接 TypeError:224-228)。这也是 advanced-chat 要做 Snapshot 的原因。
  • channel / rag-pipeline 落进 AppGenerateService.generate 都会报错app_generate_service.py:343-344)。区别是 rag-pipeline 有独立入口,channel 没有。

12. 代码地图

主题文件路径(相对克隆根)关键符号
模式分发与限流入口api/services/app_generate_service.pyAppGenerateService.generate_get_max_active_requests_build_streaming_task_on_subscribe
应用模式枚举api/models/model.pyAppMode(八个成员,:364
并发额度实现api/core/app/features/rate_limiting/rate_limit.pyRateLimit.enterRateLimit.exitRateLimitGenerator.closerate_limit_context
workflow 模式 Generatorapi/core/app/apps/workflow/app_generator.pyWorkflowAppGenerator.generate_generate_generate_worker_handle_response
advanced-chat 模式 Generatorapi/core/app/apps/advanced_chat/app_generator.pyAdvancedChatAppGenerator.generate_generate_generate_worker_handle_advanced_chat_response
agent 模式 Generatorapi/core/app/apps/agent_app/app_generator.pyAgentAppGenerator.generate_generate_worker
workflow Runnerapi/core/app/apps/workflow/app_runner.pyWorkflowAppRunner.run
advanced-chat Runnerapi/core/app/apps/advanced_chat/app_runner.pyAdvancedChatAppRunner.run
引擎事件 → 队列事件api/core/app/apps/workflow_app_runner.pyWorkflowBasedAppRunner._handle_event_publish_event
队列与停止标志api/core/app/apps/base_app_queue_manager.pyAppQueueManager.listenpublishstop_listenset_stop_flag_is_stopped_check_for_sqlalchemy_models
workflow 队列子类api/core/app/apps/workflow/app_queue_manager.pyWorkflowAppQueueManager._publish
消息类队列子类api/core/app/apps/message_based_app_queue_manager.pyMessageBasedAppQueueManager._publish
workflow 流式管线api/core/app/apps/workflow/generate_task_pipeline.pyWorkflowAppGenerateTaskPipeline.process_process_stream_response_dispatch_event_get_event_handlers
advanced-chat 流式管线api/core/app/apps/advanced_chat/generate_task_pipeline.pyAdvancedChatAppGenerateTaskPipeline.process_process_stream_responseWorkflowSnapshot
blocking/streaming 出口api/core/app/apps/base_app_generate_response_converter.pyAppGenerateResponseConverter.convert
workflow 出口详略分叉api/core/app/apps/workflow/generate_response_converter.pyconvert_stream_full_responseconvert_stream_simple_response
SSE 文本行拼装api/core/app/apps/base_app_generator.pyBaseAppGenerator.convert_to_event_stream
跨进程订阅工具api/core/app/apps/streaming_utils.pystream_topic_events
topic key 与订阅入口api/core/app/apps/message_based_app_generator.py_make_channel_keyget_response_topicretrieve_events
Celery 执行任务api/tasks/app_generate/workflow_execute_task.pyworkflow_based_app_execution_task_AppRunner.runAppExecutionParams.new_publish_streaming_response
上下文搬运api/libs/flask_utils.pypreserve_flask_contextsset_login_user
HTTP 响应收口api/libs/helper.pycompact_generate_response
停止接口api/controllers/console/app/workflow.pyWorkflowTaskStopApi.post
RAG 流水线复用api/core/app/apps/pipeline/pipeline_generator.pyapi/services/rag_pipeline/pipeline_generate_service.pyPipelineGenerator._generatePipelineGenerateService.generate