Strapi 数据迁移 WebSocket 协议解析:Remote Data Transfer 的 Dispatcher 消息模型与传输生命周期
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 |
30000 ms |
每经过该时长仍未收到匹配 uuid 的响应,则重发同一条消息(utils.ts#L84-L96 的 setInterval(sendPeriodically, retryMessageTimeout)) |
retryMessageMaxRetries |
5 |
发送计数超过该值后,以 ProviderError('error', 'Request timed out') 中止当前消息 |
retryOverrides |
无 | 单条消息级覆盖,如 assets start 使用 120000/30(remote-source#L37-L40) |
streamTimeout(pull 资产) |
300000 ms |
单个资产无进展(无新远端块、无完成写入)达到该时长即中止(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 标签,协议细节(命令集合、消息字段)以当前仓库版本的类型定义为准。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00