首页
/ Strapi 数据迁移 WebSocket 协议解析:Remote Data Transfer 的 Dispatcher 消息模型与传输生命周期

Strapi 数据迁移 WebSocket 协议解析:Remote Data Transfer 的 Dispatcher 消息模型与传输生命周期

2026-09-04 15:46:32作者:柏廷章Berta

Strapi 的远程数据迁移(Data Transfer)功能通过 WebSocket 在源端(push)或目标端(pull)Strapi 服务器之间建立一条结构化消息通道。本篇基于仓库文档 01-websocket.mdpackages/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 命令,与 initend 并列于 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 值:

  • bootstrap
  • getMetadata
  • beforeTransfer
  • getSchemas
  • rollback(仅 destination 方向)
  • close:完成一次传输(但不关闭连接)

文档提示完整且精确的消息定义见 packages/core/data-transfer/dist/strapi/remote/handlers/pull.d.tspush.d.ts——这是构建产物(dist)中的类型声明。在未构建的源码仓库中,对应的运行时实现与类型契约位于 remote/handlers/pull.tsremote/handlers/push.ts 以及协议类型目录 types/remote/protocol(其中 client/transfer/pull.tspush.tscommands.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 响应做了分类处理:

  • 401Failed to initialize the connection: Authentication Error
  • 403Failed to initialize the connection: Authorization Error
  • 404Failed 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-responsemessage(用于转发 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 顺序执行的动作:

  1. bootstrap:初始化传输环境;
  2. getMetadata:获取传输元数据;
  3. beforeTransfer:执行迁移前准备;
  4. 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: 5retryMessageTimeout: 30000(30 秒)。dispatch 内部通过 setIntervalretryMessageTimeout 毫秒重发一次同一载荷,计数超过 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 关闭阶段

清理动作

  1. 发送 close 动作:
dispatchTransferAction('close');
  1. 发送 end 命令:
dispatchCommand({ command: 'end', params: { transferID } });

连接终止

  1. 按逆序移除事件监听器:先移除 message,再 erroropenclose
  2. 关闭 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-L96setInterval(sendPeriodically, retryMessageTimeout)
retryMessageMaxRetries 5 发送计数超过该值后,以 ProviderError('error', 'Request timed out') 中止当前消息
retryOverrides 单条消息级覆盖,如 assets start 使用 120000/30remote-source#L37-L40
streamTimeout(pull 资产) 300000 ms 单个资产无进展(无新远端块、无完成写入)达到该时长即中止(remote-source#L185-L201

响应匹配机制是这套重试能够安全工作的关键:每条出站消息都会附加随机 uuidrandomUUID()),onResponse 处理器只在 response.uuid === uuid 时才解析并 resolve/reject,否则把监听重新挂回 ws.once('message', onResponse)utils.ts#L98-L132)。这保证了重发产生的重复帧、或诊断类旁路消息不会污染当前请求的响应解析。错误响应还会按 step 字段细化为不同类型的异常:transferProviderTransferErrorvalidationProviderValidationErrorinitializationProviderInitializationError

errors/providers.ts 可进一步追溯这些异常类的定义,用于在自定义客户端中做分类处理。

4. 小结:如何落地这套协议

基于文档与源码,一次自定义的远程迁移客户端应遵循如下要点:

  1. 建连:向远端 Strapi 的 /admin/transfer/runner/push(push 方向,远端为目的地)或 /admin/transfer/runner/pull(pull 方向,远端为源)建立 ws/wss 连接,并按需在 Authorization: Bearer <transfer_token> 头中携带迁移 token;url 协议只允许 http:/https:assertValidProtocol)。
  2. 顺序严格initbootstrap / getMetadata / beforeTransfer / getSchemas → 每个 step 的 start → 若干 streamendclose 动作 → end 命令 → 关闭连接;乱序会得到服务器错误。
  3. transferID 全程携带init 返回后,所有 transfer 消息自动(或手动)附带 transferID。
  4. 为慢操作放宽窗口:默认 30 秒 × 5 次重试适合大多数消息,但统计量大、数据量大的步骤(如 assets start)应使用 retryOverrides 放宽窗口,避免误报 Request timed out
  5. 流式数据要确认: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/testsremote-destination/tests 下的 checksum 协商、资产流等测试用例。需要提醒的是:该文档带有 experimental 标签,协议细节(命令集合、消息字段)以当前仓库版本的类型定义为准。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.79 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
988
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384