Strapi 数据迁移 WebSocket 协议解析:Remote Data Transfer 的 Dispatcher 消息模型与传输生命周期
【免费下载链接】strapi🚀 Strapi is the leading open-source headless CMS. It’s 100% JavaScript/TypeScript, fully customizable, and developer-first.项目地址: https://gitcode.com/GitHub_Trending/st/strapi
Strapi 的远程数据迁移(Data Transfer)功能通过 WebSocket 在源端(push)或目标端(pull)Strapi 服务器之间建立一条结构化消息通道。本篇基于仓库文档 01-websocket.md 与packages/core/data-transfer的源码实现,讲清楚三件事:WebSocket 服务端只接受哪些"传输命令"(transfer commands)以及它们必须按什么顺序发送;消息分发器(dispatcher)的dispatchCommand/dispatchTransferStep/dispatchTransferAction三个方法各自承担什么职责;以及从建连到关闭的完整传输生命周期(连接、初始化、动作、分步流式传输、关闭)如何运转,包括超时重试机制的源码级细节。读完后你将能够理解 Strapi 远程迁移协议的消息契约,并能基于 remote-source 提供者 的bootstrap()方法实现或调试自定义的 WebSocket 迁移客户端。
1. 传输命令与消息分发器
远程 WebSocket 服务器只接受特定的 WebSocket 消息——文档将其称为transfer commands。这些命令必须按特定顺序发送;如果服务器收到意外的消息,会返回错误消息。因此协议本质上是一个严格有序的"命令 → 响应"对话,而非自由的双工数据流。
文档指出,客户端应创建一个消息分发器对象(message dispatcher)来向服务器发送消息,实现位于 strapi/providers/utils.ts。阅读源码后可以确认,createDispatcher()的完整签名为:
export const createDispatcher = ( ws: WebSocket, retryMessageOptions: RetryMessageOptions = { retryMessageMaxRetries: 5, retryMessageTimeout: 30000, }, reportInfo?: (message: string) => void ) => { /* ... */ }从源码结构看,分发器的返回值正是文档所描述的三个方法,外加内部dispatch与传输状态访问器:
1.1 dispatchCommand —— 开启与结束传输
接受用于打开和关闭传输的command。源码中分发器校验的命令集合定义在 remote/handlers/constants.ts:
export const VALID_TRANSFER_COMMANDS = ['init', 'end', 'status'] as const;其中文档重点描述的两个命令:
init:初始化连接,返回transferID,此后本次传输中的所有消息都必须携带该 transferID;end:结束连接。
init命令还支持携带params。在 remote-source 提供者的initTransfer()中可以看到实际用法:当需要校验资产字节完整性时,客户端会在 init 参数中声明checksums: true,服务器若支持则回传checksums: true完成协商:
const query = this.dispatcher?.dispatchCommand({ command: 'init', ...(wantsChecksums ? { params: { transfer: 'pull', checksums: true } } : {}), });此外,从Handler接口(remote/handlers/abstract.ts)可以看到服务器侧还定义了status命令,与init、end并列于VALID_TRANSFER_COMMANDS中。
1.2 dispatchTransferStep —— 阶段切换与数据流式传输
用于在传输的阶段(step/stage)之间切换,并流式传输传输的实际数据。接受的action取值:
start:携带step值(阶段名称),表示开始该阶段;stream:可发送任意多条,携带step值与正在发送的data(例如实体数组、资产块);end:携带step值,表示该阶段结束。
在 utils.ts 的dispatchTransferStep中可以看到,stream类型的消息会要求data字段,并且所有 step 消息都会自动附加attachTransfer: true,即自动补上 transferID。
1.3 dispatchTransferAction —— 触发服务端动作
用于触发与本地提供者等价的"动作"。文档列出的 action 值:
bootstrapgetMetadatabeforeTransfergetSchemasrollback(仅 destination 方向)close:完成一次传输(但不关闭连接)
文档提示完整且精确的消息定义见packages/core/data-transfer/dist/strapi/remote/handlers/pull.d.ts与push.d.ts——这是构建产物(dist)中的类型声明。在未构建的源码仓库中,对应的运行时实现与类型契约位于 remote/handlers/pull.ts、remote/handlers/push.ts 以及协议类型目录 types/remote/protocol(其中client/transfer/pull.ts、push.ts、commands.ts分别定义了客户端消息结构),可作为阅读精确消息定义的入口。
2. 传输生命周期
原文档用一张 Mermaid 时序图完整刻画了一次传输的全过程,各阶段依次为:连接阶段 → 初始化阶段 → 传输动作阶段 → 传输步骤阶段(流式)→ 关闭阶段。下面逐阶段展开,并补充源码证据。
2.1 连接阶段:WebSocket 建连与鉴权
当 Strapi 服务器启用了数据迁移功能(即设置了admin.transfer.token.salt配置值,且server.transfer.remote.enabled未设为false)时,Strapi 会创建两个 WebSocket 服务器,路由分别为/admin/transfer/runner/pull与/admin/transfer/runner/push。源码中的路径常量印证了这一点,见 remote/constants.ts:
export const TRANSFER_PATH = '/transfer/runner' as const; export const TRANSFER_METHODS = ['push', 'pull'] as const;建立连接:在以上路由上打开 WebSocket 连接时,需要在Authorization头中提供有效的迁移 token 作为 Bearer Token:
Authorization: Bearer <transfer_token>服务器校验 token 后建立连接。文档建议参考 remote 提供者的bootstrap()方法了解初始连接的建立方式。remote-source 的bootstrap()给出了完整的建连示例:
async bootstrap(diagnostics?: IDiagnosticReporter): Promise<void> { const { url, auth } = this.options; const wsProtocol = url.protocol === 'https:' ? 'wss:' : 'ws:'; const wsUrl = `${wsProtocol}//${url.host}${trimTrailingSlash(url.pathname)}${TRANSFER_PATH}/pull`; // 未定义 auth 时,尝试公开访问迁移 if (!auth) { ws = await connectToWebsocket(wsUrl, undefined, this.#diagnostics); } // 常见的 token 鉴权,这应是主要的鉴权方式 else if (auth.type === 'token') { const headers = { Authorization: `Bearer ${auth.token}` }; ws = await connectToWebsocket(wsUrl, { headers }, this.#diagnostics); } else { throw new ProviderValidationError('Auth method not available', { check: 'auth.type', ... }); } this.ws = ws; this.dispatcher = createDispatcher(this.ws, retryMessageOptions, (message) => this.#reportInfo(message) ); const transferID = await this.initTransfer(); this.dispatcher.setTransferProperties({ id: transferID, kind: 'pull' }); await this.dispatcher.dispatchTransferAction('bootstrap'); }可以观察到几个与文档呼应的细节:HTTP(S) 地址会被转换为wss:/ws:协议后拼接/transfer/runner/pull;鉴权只支持token类型(Bearer 头),其余类型抛出ProviderValidationError;建连后立即创建 dispatcher、执行init拿到 transferID,再发出bootstrap动作——这正是文档生命周期图中"Connection → Initialization → Actions"顺序的真实代码落点。
HTTP 状态码语义:文档未展开但源码中有明确约定,connectToWebsocket 对握手阶段的非 101 响应做了分类处理:
401→Failed to initialize the connection: Authentication Error403→Failed to initialize the connection: Authorization Error404→Failed to initialize the connection: Data transfer is not enabled on the remote host- 其他状态码 →
Unexpected server response ${statusCode}
这为排查"连接失败"类问题提供了直接依据:404 意味着远端服务器未启用数据迁移功能,而不是网络不通。
事件监听器挂载:文档指出,WebSocket 创建后应立即挂载以下监听器:
'open':处理连接成功建立;'close':管理连接终止;'error':处理连接与传输错误;'message':处理来自服务器的入站消息。
在仓库实现中,connectToWebsocket内部挂载了open(resolve Promise)、unexpected-response、message(用于转发diagnostic诊断消息到诊断报告器)和error四个处理器;而close/message的完整挂载则由 dispatcher 与流式读取逻辑按阶段动态注册(见下文 2.4 的拉取流监听)。
2.2 初始化阶段
客户端发送初始命令建立传输,服务器响应唯一的transferID:
const transferID = await dispatcher.dispatchCommand('init'); // 此后所有后续消息都必须携带该 transferID从源码看,init的响应载荷除transferID外还可携带协商结果(如checksums),客户端随后调用setTransferProperties({ id: transferID, kind: 'pull' })把 transferID 与传输方向存入 dispatcher 内部状态。dispatcher 的dispatch内部逻辑会在options.attachTransfer为真时自动执行Object.assign(payload, { transferID: state.transfer?.id })(utils.ts#L58-L60),所以业务代码无需手动为每条消息附加 transferID。
2.3 传输动作阶段
通过dispatchTransferAction顺序执行的动作:
bootstrap:初始化传输环境;getMetadata:获取传输元数据;beforeTransfer:执行迁移前准备;getSchemas:获取内容类型 schemas,用于源端与目标端之间的校验。
remote-source 提供者中的getMetadata()与getSchemas()正是这两个动作的直接调用:
async getMetadata(): Promise<IMetadata | null> { const metadata = await this.dispatcher?.dispatchTransferAction<IMetadata>('getMetadata'); return metadata ?? null; }2.4 传输步骤阶段:数据流式传输
这是实际数据传输发生的主阶段,依次处理不同类型的数据(schemas、entities、assets、links、configuration):
阶段开始
dispatchTransferStep(action: "start", step)数据流式传输
dispatchTransferStep(action: "stream", step, data)阶段结束
dispatchTransferStep(action: "end", step)重试机制:数据传输期间:
- 如果在
retryMessageTimeout内未收到服务器响应; - 系统最多重试
retryMessageMaxRetries次; - 超时自动重试;
- 超过最大重试次数则中止传输。
源码给出了具体默认值与实现方式:createDispatcher的默认参数为retryMessageMaxRetries: 5、retryMessageTimeout: 30000(30 秒)。dispatch 内部通过setInterval每retryMessageTimeout毫秒重发一次同一载荷,计数超过retryMessageMaxRetries后以ProviderError('error', 'Request timed out')拒绝该 Promise。
值得注意的是仓库中的两处工程化增强,体现了该协议在实际大规模迁移中的调优思路:
- 长窗口覆盖:
dispatch支持retryOverrides选项,可对单条消息临时合并更长的重试窗口。remote-source 中定义了ASSETS_START_RETRY_OVERRIDES = { retryMessageTimeout: 120_000, retryMessageMaxRetries: 30 },因为 pull 端在收到assets阶段的start前要先跑estimateAssetTotals(数据库流式统计),大媒体库场景可能超过默认 30 秒窗口。#startStep('assets')调用时即传入该覆盖值。 - 资产停滞检测:拉取资产时,除了消息级重试,还设有独立的
streamTimeout(默认 300_000 毫秒)监测"单个资产长时间无进展",超时则销毁对应资产流(Asset ${assetID} transfer timed out)。 - 消息确认:拉取方向下,服务器推来的流式数据帧由客户端逐帧回发
{ uuid }作为确认(#respond),服务器端则用confirm()("It sends a message to the client and waits for a confirmation",见 Handler 接口)实现同一语义。这保证了"stream 一条、确认一条"的可靠流。
2.5 关闭阶段
清理动作:
- 发送 close 动作:
dispatchTransferAction('close');- 发送 end 命令:
dispatchCommand({ command: 'end', params: { transferID } });连接终止:
- 按逆序移除事件监听器:先移除
message,再error、open、close; - 关闭 WebSocket 连接。
仓库中 remote-source 的close()展示了第一步与连接关闭的结合:先dispatchTransferAction('close')完成传输,再对 ws 注册一次性close回调后调用ws.close(),连接真正关闭时才 resolve——这与文档"close 动作不关闭连接,由客户端显式关闭"的描述一致。
3. 消息超时与重试模型
原文档用一张状态图描述了消息-响应协议的状态迁移:Init → Ready → Dispatching → WaitingResponse → (Retrying | Ready) → Error,并说明:因为传输依赖"消息→响应"协议,如果 WebSocket 服务器无法回复(例如网络不稳定),连接就会停摆。为此,每个提供者的选项中都包含retryMessageOptions:在达到给定超时后重发消息,并在给定次数的失败重试后中止传输。
结合源码,该模型可以细化为以下可验证事实:
| 配置项 | 默认值 | 行为(源码依据) |
|---|---|---|
retryMessageTimeout | 30000ms | 每经过该时长仍未收到匹配 uuid 的响应,则重发同一条消息(utils.ts#L84-L96 的setInterval(sendPeriodically, retryMessageTimeout)) |
retryMessageMaxRetries | 5 | 发送计数超过该值后,以ProviderError('error', 'Request timed out')中止当前消息 |
retryOverrides | 无 | 单条消息级覆盖,如 assetsstart使用120000/30(remote-source#L37-L40) |
streamTimeout(pull 资产) | 300000ms | 单个资产无进展(无新远端块、无完成写入)达到该时长即中止(remote-source#L185-L201) |
响应匹配机制是这套重试能够安全工作的关键:每条出站消息都会附加随机uuid(randomUUID()),onResponse处理器只在response.uuid === uuid时才解析并 resolve/reject,否则把监听重新挂回ws.once('message', onResponse)(utils.ts#L98-L132)。这保证了重发产生的重复帧、或诊断类旁路消息不会污染当前请求的响应解析。错误响应还会按step字段细化为不同类型的异常:transfer→ProviderTransferError、validation→ProviderValidationError、initialization→ProviderInitializationError。
从 errors/providers.ts 可进一步追溯这些异常类的定义,用于在自定义客户端中做分类处理。
4. 小结:如何落地这套协议
基于文档与源码,一次自定义的远程迁移客户端应遵循如下要点:
- 建连:向远端 Strapi 的
/admin/transfer/runner/push(push 方向,远端为目的地)或/admin/transfer/runner/pull(pull 方向,远端为源)建立ws/wss连接,并按需在Authorization: Bearer <transfer_token>头中携带迁移 token;url协议只允许http:/https:(assertValidProtocol)。 - 顺序严格:
init→bootstrap/getMetadata/beforeTransfer/getSchemas→ 每个 step 的start→ 若干stream→end→close动作 →end命令 → 关闭连接;乱序会得到服务器错误。 - transferID 全程携带:
init返回后,所有 transfer 消息自动(或手动)附带 transferID。 - 为慢操作放宽窗口:默认 30 秒 × 5 次重试适合大多数消息,但统计量大、数据量大的步骤(如 assets
start)应使用retryOverrides放宽窗口,避免误报Request timed out。 - 流式数据要确认:pull 方向下服务器推送的流式帧需要以
{ uuid }回应,形成可靠的逐帧确认闭环。
本文全部内容以 docs/docs/docs/01-core/data-transfer/02-providers/05-remote-strapi/01-websocket.md 的协议描述为骨架,并以 packages/core/data-transfer 中的分发器实现(strapi/providers/utils.ts)、pull/push 处理器(strapi/remote/handlers)与 remote-source 提供者(strapi/providers/remote-source/index.ts)作为实现级佐证;相关行为还可参考 packages/core/data-transfer/src/strapi/providers/remote-source/tests与 remote-destination/tests下的 checksum 协商、资产流等测试用例。需要提醒的是:该文档带有experimental标签,协议细节(命令集合、消息字段)以当前仓库版本的类型定义为准。
【免费下载链接】strapi🚀 Strapi is the leading open-source headless CMS. It’s 100% JavaScript/TypeScript, fully customizable, and developer-first.项目地址: https://gitcode.com/GitHub_Trending/st/strapi
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考