SFTP 多目标文件分发设计
引言
需求背景
JumpServer Client 在 SFTP 资产连接里,需要一套文件传输能力。最初只做两件事:两个面板之间的拖拽,以及面板之间的单文件传输。后续才扩成多文件批量选择、多目标机器分发、同名文件处理、独立进度展示,以及页面刷新后的恢复
功能确定后,首先要考虑的下面这几个贯穿整套设计的核心问题
- 一次选择 3 个文件并发送到 2 台机器时,系统应该如何拆分任务,又应该以什么层级向用户展示
- 当其中一台目标机器离线或传输失败时,其他目标机器是否应该继续传输,如何实现目标之间的故障隔离
- 目标端存在同名文件时,覆盖、跳过或保留两者的策略应该由客户端决定,还是由目标端执行
- 页面显示传输进度为 80% 时,这 80% 代表浏览器已经发送的数据量,还是目标端已经确认写入的数据量
- 文件传输完成后,如何校验目标端最终生成的文件与源文件内容完全一致
- 页面刷新、连接中断或目标端写入确认丢失后,客户端应该从哪个字节位置继续,如何证明这个断点是安全的
如果这些问题继续堆在 Vue 组件和拖拽回调里,组件会同时承担 UI、连接管理、协议通信、任务调度、持久化和校验职责,不仅难以维护,也很难证明续传正确性
因此,改造的第一步是把 "用户操作" 转换为可被调度、持久化和恢复的任务模型:将一次分发操作,分成三级
任务模型
三级任务结构与故障隔离
假设用户选择 2 个文件,想要发送到 2 台机器上,那么系统就会创建 2 个目标的 batch 和 4 个文件的 task
一次分发操作
├── 目标机器 A 的 batch
│ ├── 文件 task 1
│ └── 文件 task 2
└── 目标机器 B 的 batch
├── 文件 task 1
└── 文件 task 2| 层级 | 对应含义 | 用途 |
|---|---|---|
| 分发 | 用户点一次 "开始分发" | 只存在于认知和 batchId 前缀里,不单独落盘 |
batch | 一台目标机器 | 故障隔离和冲突处理的边界,这台挂了只停这一组 |
task | 一个文件传到一台机器 | 真正传输、暂停、重试、取消、记进度的最小单元 |
最上层的分发操作服务于用户认知,用户会把这次点击 "开始分发" 视为一个整体
持久化层只保存 batches 和 tasks,分发操作不额外创建第三种实体,而是通过 batch ID 中共同的 distribution ID 在 UI 中动态聚合
这套三级结构主要存在于客户端调度与展示层,后端不会感知 "一次分发操作" 或 batch,它只处理一个 transferId 对应的一次文件传输,这让后端协议保持单一职责,也让前端可以独立调整聚合方式
三级转换流程
用户点一次 "开始分发"
这时系统只知道四件事:从哪来(
sourceEndpoint/sourcePath)、要传哪些文件(entries)、传到哪几台机器(targets)、碰上同名文件怎么办(conflictPolicy)按目标和文件拆开,再排队
一台机器上的每个文件各算一条任务。同一台机器的任务归成一组。组进了队列之后,才真正有批次和单条任务(
enqueueBatch)展示传输效果时,再按一次分发拼回去
底下存的是一条条任务。展示时先认出哪些任务是同一次点出来的,再按目标机器分开,才能画出 "一次分发 → 多台机器 → 多个文件"(
groupSftpTransferBatches)相关代码实现
接口定义
type FileTransferConflictPolicy = "ask" | "overwrite" | "skip" | "keep_both" interface FileTransferEndpointRef { /** 一次授权会话里稳定,刷新后可能变 */ id: string label: string } interface FileTransferSource { path: string name: string size: number } interface SftpDistributionEntry { name: string /** 列表里可能是字符串,展开时会先 `Number()` */ size: string | number } interface SftpDistributionTarget { endpoint: FileTransferEndpointRef destinationPath: string } interface CreateFileTransferTaskInput { /** 同一目标共用;有分发时形如 `sftp-dist:{id}::target:{endpointId}` */ batchId: string sourceEndpoint: FileTransferEndpointRef destinationEndpoint: FileTransferEndpointRef source: FileTransferSource destinationPath: string conflictPolicy: FileTransferConflictPolicy } interface SftpDistributionInput { /** 这一次点击的 id,只用来拼 `batchId` 前缀,不单独落盘 */ distributionId?: string /** 源面板 */ sourceEndpoint: FileTransferEndpointRef /** 源目录,例如 `/home/data` */ sourcePath: string /** 选中的文件 */ entries: SftpDistributionEntry[] /** 目标机器和各自的目标目录 */ targets: SftpDistributionTarget[] /** 同名策略,整次分发共用 */ conflictPolicy: FileTransferConflictPolicy } interface SftpDistributionGroup { destination: FileTransferEndpointRef /** 这台目标上要入队的全部 task */ inputs: CreateFileTransferTaskInput[] }具体实现
sftpDistribution.ts/** * 把一次分发展开成按目标切开的任务组 * * @param input 源面板、选中文件、目标机器和冲突策略 * @returns 每个目标一组,`inputs` 交给 `enqueueBatch` */ export function buildSftpDistributionGroups( input: SftpDistributionInput ): SftpDistributionGroup[] { const sourcePath = normalizeDirectory(input.sourcePath) /** * 文件名必须有,`size` 必须是 `>= 0` 的有限数字 * 列表里的 `size` 可能是字符串,先转成 number */ const entries = input.entries .map((entry) => ({ ...entry, size: Number(entry.size) })) .filter((entry) => entry.name && Number.isFinite(entry.size) && entry.size >= 0) return input.targets /** 源和目标是同一个 endpoint 时丢掉,避免自己传给自己 */ .filter((target) => target.endpoint.id !== input.sourceEndpoint.id) .map((target) => { const destinationPath = normalizeDirectory(target.destinationPath) return { destination: target.endpoint, /** 这台目标 × 每个文件 = 一组 `CreateFileTransferTaskInput` */ inputs: entries.map((entry) => ({ /** * 有 `distributionId` 时,同一目标共用一个 `batchId` * 前缀相同,后面按目标切开:`sftp-dist:{id}::target:{endpointId}` */ batchId: input.distributionId ? `sftp-dist:${input.distributionId}::target::${encodeURIComponent(target.endpoint.id)}` : "", sourceEndpoint: input.sourceEndpoint, destinationEndpoint: target.endpoint, source: { name: entry.name, size: entry.size, path: `${sourcePath}/${entry.name}`.replace(/\/+/g, "/") }, destinationPath, conflictPolicy: input.conflictPolicy })) } }) /** 文件全被滤掉时,这台目标不要留空组 */ .filter((group) => group.inputs.length > 0) }useSftpTransferCenterSelectors.ts/** * `batchId` = `sftp-dist:{distributionId}::target:{endpointId}` * 截掉 `::target::` 后面,同一拨分发就回到同一个 `groupId` * * @param batchId 单个目标 batch 的 id * @returns 分发层的聚合 id;对不上分隔符时原样返回 */ function toBatchGroupId(batchId: string): string { const index = batchId.indexOf("::target::") return index >= 0 ? batchId.slice(0, index) : batchId } /** * 展示时把扁平的 task 列表重新聚成「一次分发 → 多个目标」 * * @param allTasks 传输中心里当前能看到的全部 task * @returns 按分发 id 聚合,组内再按目标 endpoint 切开 */ function groupSftpTransferBatches( allTasks: FileTransferTask[] ): SftpTransferBatchGroup[] { const groups = new Map<string, FileTransferTask[]>() for (const task of allTasks) { const groupId = toBatchGroupId(task.batchId) const grouped = groups.get(groupId) ?? [] grouped.push(task) groups.set(groupId, grouped) } return [...groups.entries()].map(([id, batchTasks]) => { /** 同一拨分发里,再按目标 endpoint 切开 */ const targetMap = new Map<string, FileTransferTask[]>() for (const task of batchTasks) { const targetTasks = targetMap.get(task.destinationEndpoint.id) ?? [] targetTasks.push(task) targetMap.set(task.destinationEndpoint.id, targetTasks) } return { id, tasks: batchTasks, targets: [...targetMap.entries()].map(([endpointId, tasks]) => ({ endpointId, tasks })) } }) }webSftp_transfer.go// 后端只看见单个 transferId,也就是前端那个 task.id // 没有 distribution,也没有 batch func (h *webSftp) handleTransferMutation( request *webSftpRequest, msg *Message, response *Message, ) { switch msg.Cmd { case "transfer_prepare": // request.TransferID / Path / Size / ConflictPolicy result, err = h.volume.prepareTransfer( request.TransferID, request.Path, request.Size, request.ConflictPolicy, ) case "transfer_write": // 再带上 offset、sha256 和这一片的字节 result, err = h.volume.writeTransferChunk( request.TransferID, request.Path, request.Size, request.OffSet, request.SHA256, msg.Raw, ) case "transfer_commit": result, err = h.volume.commitTransfer( request.TransferID, request.Path, request.Size, request.SHA256, request.ConflictPolicy, ) } }
任务状态机管理任务生命周期
文件任务共有 9 个状态
queued:任务已进入调度队列,等待源端和目标端可用preparing:与目标端协商冲突策略、临时文件和可续传断点transferring:按 2 MB 分片读取、校验并写入verifying:全部分片已写入,正在生成完整 SHA-256 并提交paused:暂时停止并保留断点,等待连接恢复或用户继续completed:目标端已经完成最终提交skipped:依据冲突策略跳过,属于正常终态failed:校验、协议或提交发生不可自动恢复的错误canceled:用户主动取消,系统会尝试清理目标端临时状态
正常路径为 queued → preparing → transferring → verifying → completed
连接断开、WebSocket 不可用或 endpoint 尚未重新注册时进入 paused,因为这类问题通常是暂时不可用,已经确认的断点仍有价值;数据哈希不一致、响应格式非法或 commit 失败时进入 failed,需要用户明确重试
相关代码
/** 前端任务生命周期,调度、展示、持久化都看这个 */
type FileTransferStatus =
| "queued"
| "preparing"
| "transferring"
| "verifying"
| "paused"
| "completed"
| "skipped"
| "failed"
| "canceled"
/**
* 目标端 prepare / status 的回包
* `state` 只有协议协商结果,不是上面那 9 个任务状态
*/
interface FileTransferResumeState {
transferId: string
committedBytes: number
totalBytes: number
state: "ready" | "completed" | "skipped" | "missing" | "conflict"
}/** 刷新后还能接着跑的中间态,恢复时一律先改成 `paused` */
const resumableStatuses = new Set<FileTransferStatus>([
"queued",
"preparing",
"transferring",
"verifying"
])
/** 已经结束,不再进调度队列 */
const terminalStatuses = new Set<FileTransferStatus>([
"completed",
"skipped",
"failed",
"canceled"
])
/**
* 从 IndexedDB 读回任务后,中间态先停住
* 连接和授权还没恢复,不能自动接着传
*/
tasks.value = persisted.tasks.map((task) =>
resumableStatuses.has(task.status)
? { ...task, status: "paused", updatedAt: Date.now() }
: task
)/**
* 把目标端的协议状态映射成任务状态
* @param prepared 这次 `prepareTransfer` 的回包
*/
if (prepared.state === "skipped") {
patchTask(task.id, { status: "skipped", confirmedBytes: prepared.committedBytes })
return
}
/** 同名冲突:整组 batch 一起停,等用户选覆盖 / 跳过 / 都保留 */
if (prepared.state === "conflict") {
for (const batchTask of tasks.value) {
if (batchTask.batchId === task.batchId && !terminalStatuses.has(batchTask.status)) {
patchTask(batchTask.id, { status: "paused", error: conflictError })
}
}
return
}
if (prepared.state === "completed") {
patchTask(task.id, { status: "completed", confirmedBytes: prepared.totalBytes })
return
}
/** `ready` 才能往下走分片传输 */
if (prepared.state !== "ready" || prepared.totalBytes !== task.source.size) {
throw new Error("Invalid file transfer preparation response")
}paused 与 canceled 的区别尤其重要:暂停是可恢复的中间态,取消是主动放弃后的终态;skipped 虽然没有传输文件内容,但已经按用户策略完成处理,因此也必须计入批次的完成统计,不能让聚合进度永远停留在 loading 状态
状态数量不是为了让 UI 显得复杂,而是为了让调度器能够根据语义决定下一步动作;后端只返回与协议协商相关的 ready、conflict、skipped 和 completed,客户端再将协议结果映射为完整任务生命周期
后端代码
// 协议回包。State 只有 ready / completed / skipped / missing / conflict
type webSftpTransferResponse struct {
TransferID string `json:"transfer_id"`
CommittedBytes int64 `json:"committed_bytes"`
TotalBytes int64 `json:"total_bytes"`
State string `json:"state"`
Duplicate bool `json:"duplicate,omitempty"`
}
// 目标路径已经存在时,按冲突策略回协商状态
if _, statErr := u.UserSftp.Stat(targetPath); statErr == nil {
if conflictPolicy == "skip" {
return webSftpTransferResponse{
TransferID: transferID,
TotalBytes: totalSize,
State: "skipped",
}, nil
}
if conflictPolicy == "ask" {
return webSftpTransferResponse{
TransferID: transferID,
TotalBytes: totalSize,
State: "conflict",
}, nil
}
}后端只回 ready、conflict、skipped、completed 这几个协商结果,9 个任务状态是前端自己管的。刷新后还没跑完的任务一律先停住,等连接回来再继续
进度与校验
可信进度条
浏览器把分片发出去,只能说明数据离开了本机,不能说明后端已经写进文件。如果这时就把进度推到 80%,断网、ACK 丢了、对面只写了一半,页面上仍可能停在 80% 甚至 100%。刷新后再按这个数续传,就会从错误的位置继续拼接文件
在当前这套设计中,进度条是根据 "目标端已经确认写入的字节数" 来计算的。目标端完成写入后会返回 committedBytes 和 duplicate,其中 committedBytes 表示目标端已经确认写入到哪个字节位置,duplicate 表示本次是否为重复分片
全部流程为
- 源端读取分片
- Worker 校验分片 SHA-256
- 目标端写入分片
- 目标端返回 ACK
- 更新
confirmedBytes - 持久化断点与哈希状态
重要
客户端只有收到 ACK,才会把服务端的 committedBytes 更新为本地任务的 confirmedBytes,并将这个进度与对应的 checksumState 一起保
接口定义
/** 发给目标端的一片数据 */
interface FileTransferWriteInput {
transferId: string
targetPath: string
totalBytes: number
offset: number
data: Uint8Array
sha256: string
}
/** 目标端对这一片的确认 */
interface FileTransferWriteAck {
/** 目标端已经写到哪个字节,进度只信这个 */
committedBytes: number
/** 这一片之前已经写过,并且内容和哈希对得上 */
duplicate: boolean
}具体实现
/**
* 从 WebSocket 回包里取出 ACK
* `committed_bytes` 必须是安全整数,否则整片作废
*/
export function parseSftpTransferWriteAck(
message: SftpDataMessage
): FileTransferWriteAck {
const response = parseTransferPayload<TransferResponseWire>(message)
const committedBytes = response.committed_bytes
if (
typeof committedBytes !== "number" ||
!Number.isSafeInteger(committedBytes)
) {
throw new TypeError("Invalid SFTP transfer write acknowledgement")
}
return {
committedBytes,
duplicate: Boolean(response.duplicate)
}
}/**
* 发出去不算数,ACK 过了边界检查才推进本地进度
* `confirmedBytes` 和 `checksumState` 必须一起存,后面续传才对得上
*/
const ack = await destination.writeChunk({
transferId: task.id,
targetPath: currentTargetPath,
totalBytes: task.source.size,
offset,
data: chunk.data,
sha256: checksum.chunkChecksum
})
if (
ack.committedBytes < offset + chunk.data.length ||
ack.committedBytes > task.source.size
) {
throw new Error("Invalid file transfer write acknowledgement")
}
offset = ack.committedBytes
checksumState = checksum.state
patchTask(task.id, { confirmedBytes: offset, checksumState })// offset 超过已写入位置:乱序,直接拒绝
if committedBytes > totalSize || offset > committedBytes {
return webSftpTransferResponse{},
fmt.Errorf("file transfer chunk offset is out of order")
}
// offset 落在已写入区间里:当重复片,读出来对哈希
// 对得上就回 duplicate,进度停在原来的 committedBytes,不往前加
if offset < committedBytes {
existing := make([]byte, len(data))
n, readErr := file.ReadAt(existing, offset)
if readErr != nil && readErr != io.EOF {
return webSftpTransferResponse{}, readErr
}
if n != len(data) ||
!strings.EqualFold(sha256Hex(existing), expectedSHA256) {
return webSftpTransferResponse{},
fmt.Errorf("file transfer duplicate chunk does not match")
}
return webSftpTransferResponse{
TransferID: transferID,
CommittedBytes: committedBytes,
TotalBytes: totalSize,
State: "ready",
Duplicate: true,
}, nil
}
// offset 正好接在末尾:写入,ACK 带回新的 committedBytes
if _, err = file.WriteAt(data, offset); err != nil {
return webSftpTransferResponse{}, err
}
return webSftpTransferResponse{
TransferID: transferID,
CommittedBytes: offset + int64(len(data)),
TotalBytes: totalSize,
State: "ready",
}, nil前端刻意没有在 writeChunk 发出后立即更新进度,而是先校验 ACK 的边界;后端对重复 offset 读取已有内容并比较 SHA-256,只有内容一致才返回 duplicate: true,因此 ACK 同时承担进度确认与幂等写入证明
因此,界面展示的是 "目标端已确认落盘的进度",而不是 "客户端已经发送的进度",这可以避免网络中断时出现看似 100% 但目标文件不完整的假进度。如果分片已经发出,但 ACK 返回前连接断开,客户端不会自行推断分片成功,而是在恢复时重新询问目标端的实际断点
双端对账决策
双端对账发生在任务每次真正开始执行时,包括首次传输、失败重试和暂停恢复
任务进入 preparing 后,客户端携带原来的 transferId、目标路径、文件名、文件大小和冲突策略调用目标端的 prepareTransfer,目标端返回当前状态及已经确认写入的 committedBytes
客户端随后对比目标端的 committedBytes、本地的 confirmedBytes 和 checksumState
- 两端字节数一致,且非零断点拥有对应的增量哈希状态时,从该 offset 继续
- 两端字节数不一致时,无法证明本地哈希现场与目标端数据对应,从零重传
- 两端字节数一致但本地缺少增量哈希状态时,仍然无法安全计算最终文件哈希,从零重传
从零重传前,客户端会通知目标端丢弃旧的临时传输状态,再重新执行 prepare,并确认断点已经回到零
这套规则可以概括为 "能证明一致才续传,否则安全重传",它允许在证据不足时牺牲一部分网络成本,但不会用不可信的断点拼接文件
双端对账之所以放在每次 preparing,而不是只在页面恢复时执行,是因为 ACK 丢失、重试和临时连接断开都可能让两端检查点短暂分叉;后端以 .part 文件真实大小作为事实,前端只有在字节数和哈希现场同时对应时才接受这个断点
查看续传与安全重传的判断
// ui/store/modules/fileTransfer.ts
let prepared = await destination.prepareTransfer(prepareInput);
const checksumAligned =
prepared.committedBytes === task.confirmedBytes &&
(prepared.committedBytes === 0 || Boolean(task.checksumState));
if (!checksumAligned) {
await destination.cancelTransfer({
transferId: task.id,
targetPath: currentTargetPath,
discard: true
});
prepared = await destination.prepareTransfer(prepareInput);
if (
prepared.state !== "ready" ||
prepared.committedBytes !== 0 ||
prepared.totalBytes !== task.source.size
) {
throw new Error("Unable to restart inconsistent file transfer state");
}
}
let offset = prepared.committedBytes;
let checksumState = offset === 0 ? "" : task.checksumState;// pkg/httpd/websftp_transfer.go
func (u *UserWebVolume) prepareTransfer(
transferID, targetPath string,
totalSize int64,
conflictPolicy string,
) (webSftpTransferResponse, error) {
stagePath, err := transferStagePath(targetPath, transferID)
if err != nil {
return webSftpTransferResponse{}, err
}
if info, statErr := u.UserSftp.Stat(stagePath); statErr == nil {
if info.Size() > totalSize {
return webSftpTransferResponse{},
fmt.Errorf("file transfer stage exceeds expected size")
}
return webSftpTransferResponse{
TransferID: transferID,
CommittedBytes: info.Size(),
TotalBytes: totalSize,
State: "ready",
}, nil
}
file, err := u.UserSftp.Create(stagePath)
if err != nil {
return webSftpTransferResponse{}, err
}
_ = file.Close()
return webSftpTransferResponse{
TransferID: transferID,
TotalBytes: totalSize,
State: "ready",
}, nil
}IndexedDB 检查点
IndexedDB 数据库名为 jumpserver-file-transfer,对象仓库名为 tasks,固定 record key 为 state
系统没有为每个任务分别创建 IndexedDB key,而是在 state 下保存一份完整快照
type PersistedTransferState = {
batches: TransferBatch[]
tasks: TransferTask[]
}batch 主要保存批次 ID、任务 ID 列表和创建时间;task 保存任务 UUID、batch ID、源端与目标端引用、源文件路径、文件名、大小、目标目录、冲突策略、当前状态、confirmedBytes、checksumState、最终 checksum、错误信息及时间戳
IndexedDB 不保存文件二进制、SFTP 凭证、WebSocket 或 Worker 实例,它保存的是 "任务控制信息 + 可恢复检查点"
每次任务状态变化或收到分片 ACK 后,客户端都会更新内存中的任务,并异步安排持久化;连续保存通过 microtask 合并,避免每个细小状态变化都立即触发一次 IndexedDB 写入
页面刷新后,系统读取这份快照,并把刷新前处于 queued、preparing、transferring 或 verifying 的任务统一恢复为 paused,避免连接和授权尚未恢复时自动在后台发起传输
客户端与目标端各保存一半恢复证据:IndexedDB 保存任务语义和哈希现场,目标端隐藏临时文件保存真实字节;任何一端都不能单独证明续传安全
查看双端检查点分别保存在哪里
// ui/shared/file-transfer/persistence.ts
const databaseName = "jumpserver-file-transfer";
const storeName = "tasks";
const recordKey = "state";
function createPersistedSnapshot(state: PersistedFileTransferState) {
return {
batches: state.batches.map((batch) => ({
id: batch.id,
taskIds: [...batch.taskIds],
createdAt: batch.createdAt
})),
tasks: state.tasks.map((task) => ({
id: task.id,
batchId: task.batchId,
sourceEndpoint: { ...task.sourceEndpoint },
destinationEndpoint: { ...task.destinationEndpoint },
source: { ...task.source },
destinationPath: task.destinationPath,
conflictPolicy: task.conflictPolicy,
status: task.status,
confirmedBytes: task.confirmedBytes,
checksumState: task.checksumState,
checksum: task.checksum,
error: task.error,
createdAt: task.createdAt,
updatedAt: task.updatedAt
}))
};
}
objectStore.put(createPersistedSnapshot(state), recordKey);// pkg/httpd/websftp_transfer.go
func transferStagePath(targetPath, transferID string) (string, error) {
if targetPath == "" || transferID == "" || len(transferID) > 128 {
return "", fmt.Errorf("invalid file transfer request")
}
base := path.Base(targetPath)
return path.Join(
path.Dir(targetPath),
fmt.Sprintf(".%s.jms-transfer-%s.part", base, transferID),
), nil
}
func (u *UserWebVolume) cancelTransfer(
transferID, targetPath string,
discard bool,
) (webSftpTransferResponse, error) {
stagePath, err := transferStagePath(targetPath, transferID)
if discard {
err = u.UserSftp.Remove(stagePath)
}
return webSftpTransferResponse{
TransferID: transferID,
State: "ready",
}, err
}增量 SHA-256 状态
只保存 confirmedBytes 可以知道下一次从哪里读,但无法从断点继续计算整个文件的 SHA-256
SHA-256 不能通过 "前半段哈希 + 后半段哈希" 合并成完整文件哈希,如果一个 4 GB 文件已经传输了 3 GB,却只保存了字节位置,那么恢复后为了得到最终哈希,仍然需要重新读取并计算前面的 3 GB
因此,Worker 会把哈希计算现场序列化为 checksumState
type IncrementalSha256State = {
h: number[]
buffer: number[]
length: number
}h 是 SHA-256 内部的 8 个中间哈希值,buffer 保存尚未组成一个 64 字节哈希块的尾部数据,length 表示已经参与计算的总字节数
每处理一个分片,Worker 都会在旧状态上继续计算并返回新状态;主线程只有在收到目标端 ACK 后,才把新的 confirmedBytes 和 checksumState 作为同一个可信检查点保存
刷新后 Worker 虽然被销毁,但新的 Worker 可以反序列化这份状态,从断点继续计算,不需要缓存文件内容,也不需要为了最终校验而重复读取已经确认的部分
Worker 校验
Worker 承担两类 SHA-256 计算
第一类是分片校验,源端读取分片时会给出分片 SHA-256,Worker 对浏览器实际收到的字节重新计算,两者不一致就立即停止任务,从而发现分片错位、截断、编码异常或传输损坏
第二类是完整文件校验,Worker 持续更新增量 SHA-256,在所有分片写入完成后生成最终 checksum,并在 commit 阶段交给目标端进行最终完整性确认
哈希属于 CPU 密集工作,把它放在 Worker 中可以避免大文件连续计算阻塞 Vue 主线程,减少进度动画、菜单和传输中心交互的卡顿
实现中会先复制一份分片,再通过 transferable 把副本交给 Worker,因为 ArrayBuffer 被转移后会在主线程中 detach,而原始分片随后仍要发送给目标端
分片校验与完整文件校验并不是重复工作:分片 hash 用于尽早发现当前块损坏,全文件 hash 用于确认所有分片按正确顺序组成了最终文件;后端 commit 还会重新读取 .part 文件计算一次完整 SHA-256,形成端到端校验
查看 Worker 与后端完整性校验
// ui/workers/fileTransferChecksum.worker.ts
interface Sha256State {
h: number[];
buffer: number[];
length: number;
}
function parseState(value: string): Sha256State {
if (!value) {
return { h: [...initialHash], buffer: [], length: 0 };
}
const state = JSON.parse(value) as Sha256State;
if (
state.h.length !== 8 ||
state.buffer.length >= 64 ||
!Number.isSafeInteger(state.length)
) {
throw new Error("Invalid file transfer checksum state");
}
return state;
}
workerScope.onmessage = async (event) => {
const request = event.data;
const state = parseState(request.state);
if (request.kind === "finalize") {
workerScope.postMessage({
id: request.id,
checksum: finalize(state),
state: JSON.stringify(state)
});
return;
}
const bytes = new Uint8Array(request.data);
const chunkChecksum = hex(
new Uint8Array(await crypto.subtle.digest("SHA-256", bytes))
);
update(state, bytes);
workerScope.postMessage({
id: request.id,
chunkChecksum,
state: JSON.stringify(state)
});
};// ui/shared/file-transfer/checksum.ts
export async function updateFileTransferChecksum(
state: string,
bytes: Uint8Array
) {
const copy = bytes.slice();
const response = await send(
{ kind: "update", state, data: copy.buffer },
[copy.buffer]
);
return {
chunkChecksum: response.chunkChecksum,
state: response.state
};
}// pkg/httpd/websftp_transfer.go
file, err := u.UserSftp.Open(stagePath)
if err != nil {
return webSftpTransferResponse{}, err
}
info, err := file.Stat()
if err != nil || info.Size() != totalSize {
return webSftpTransferResponse{},
fmt.Errorf("file transfer is incomplete")
}
hash := sha256.New()
if _, err = io.Copy(hash, file); err != nil {
return webSftpTransferResponse{}, err
}
if !strings.EqualFold(
hex.EncodeToString(hash.Sum(nil)),
expectedSHA256,
) {
return webSftpTransferResponse{},
fmt.Errorf("file transfer checksum mismatch")
}冲突与调度
冲突策略为什么由目标端执行
默认 ask 策略在 prepare 阶段遇到同名文件时,会暂停对应目标 batch,等待用户选择 overwrite、skip 或 keep_both
overwrite:任务重新排队,由目标端执行替换skip:冲突任务直接进入skipped,不读取也不发送文件内容keep_both:客户端传递策略,由目标端生成不冲突的新文件名
具体命名和原子替换放在目标端,可以避免 Web、Tauri 和其他客户端各自实现不同规则,也能更好地处理多个客户端同时写入时的文件名竞态
后端当前按照 名称 (序号).扩展名 生成候选路径,例如 report (1).pdf;没有独立扩展名的点文件会保留完整名称,例如 .env 会变为 .env (1)
候选路径选择与最终重命名在当前进程内由互斥锁串行执行,如果其他进程在 Stat 与 Rename 之间占用了候选路径,后端会识别冲突并继续尝试下一个序号
之所以在前端按目标 batch 统一解决冲突,是为了避免同一台机器上的多个文件分别打断用户;具体文件名选择和重命名必须留在后端,因为只有后端能在真正落盘时观察并发竞争
查看冲突策略如何贯穿前后端
// ui/store/modules/fileTransfer.ts
function resolveBatchConflict(
batchId: string,
conflictPolicy: Exclude<FileTransferConflictPolicy, "ask">
) {
for (const task of tasks.value) {
if (
task.batchId !== batchId ||
terminalStatuses.has(task.status)
) {
continue;
}
patchTask(task.id, {
conflictPolicy,
status: "queued",
error: undefined
});
}
kick();
}// pkg/httpd/websftp_transfer.go
filename := path.Base(targetPath)
extension := path.Ext(filename)
// .env 没有独立扩展名,结果为 .env (1)
if strings.TrimSuffix(filename, extension) == "" {
extension = ""
}
base := strings.TrimSuffix(filename, extension)
for index := 1; index <= 10000; index++ {
candidate := path.Join(
path.Dir(targetPath),
fmt.Sprintf("%s (%d)%s", base, index, extension),
)
occupied, err := exists(candidate)
if err != nil {
return "", err
}
if !occupied {
return candidate, nil
}
}// pkg/httpd/websftp_transfer.go
switch conflictPolicy {
case "keep_both":
targetPath, err = commitKeepBothTarget(
targetPath,
u.transferTargetExists,
func(candidate string) error {
return u.UserSftp.Rename(stagePath, candidate)
},
)
case "overwrite":
err = u.UserSftp.PosixRename(stagePath, targetPath)
default:
err = u.UserSftp.Rename(stagePath, targetPath)
}endpoint 级互斥形成天然背压
调度器启动一个任务前,会同时检查源 endpoint 和目标 endpoint 是否空闲,只有两端都没有被其他任务占用时才开始传输
这意味着同一条 SFTP 或 WebSocket 连接不会被多个文件任务同时读写,互不共享 endpoint 的任务仍然可以并行;同一个源向多台机器分发时,因为任务共享源 endpoint,当前实现更偏向串行,而不是无限并发广播
它牺牲了部分峰值吞吐,换来了以下收益
- 避免同一连接上的消息交错和状态竞争
- 控制浏览器内存、Worker 计算和服务端写入压力
- 让连接中断、暂停和恢复的边界更容易推理
- 通过任务排队自然形成背压,不需要再维护一套复杂的全局并发计数
文件本身按 2 MB 分片处理,也让内存占用与文件总大小解耦
这项取舍不是 "系统不支持并发",而是 "并发边界设在独立 endpoint 之间",优先保证正确性和连接稳定性,再为后续的连接池或源端多路读取留下演进空间
这里的互斥不是 UI 层禁用按钮,而是调度器的资源占用模型;后端再用 2 MB 上限和 offset 顺序校验约束单次请求,前后端共同控制连接、内存与落盘压力
查看背压如何落到调度与协议校验
// ui/store/modules/fileTransfer.ts
function kick() {
for (const queue of destinationQueues.value.values()) {
const next = queue.find(
(task) =>
task.status === "queued" &&
!runningEndpoints.has(task.sourceEndpoint.id) &&
!runningEndpoints.has(task.destinationEndpoint.id)
);
if (!next) continue;
runningEndpoints.add(next.sourceEndpoint.id);
runningEndpoints.add(next.destinationEndpoint.id);
const execution = runTask(next.id);
runningTasks.set(next.id, execution);
void execution.finally(() => {
runningTasks.delete(next.id);
runningEndpoints.delete(next.sourceEndpoint.id);
runningEndpoints.delete(next.destinationEndpoint.id);
kick();
});
}
}// pkg/httpd/websftp_transfer.go
const transferChunkMaxSize = 2 * 1024 * 1024
if offset < 0 ||
len(data) == 0 ||
offset+int64(len(data)) > totalSize ||
expectedSHA256 != sha256Hex(data) {
return webSftpTransferResponse{},
fmt.Errorf("invalid file transfer chunk")
}
if committedBytes > totalSize || offset > committedBytes {
return webSftpTransferResponse{},
fmt.Errorf("file transfer chunk offset is out of order")
}恢复与演进
页面刷新恢复的能力与边界
页面刷新后,IndexedDB 可以恢复 batch、task、confirmedBytes 和 checksumState,但它不能恢复 SFTP 连接、WebSocket、授权状态或 Worker 实例
当原源端和目标端重新注册为可用的 FileTransferEndpoint 后,用户可以继续任务,系统使用原 transferId 重新 prepare,并通过双端对账决定续传或重传
当前多远端场景仍有一项明确边界:endpoint ID 与会话 token 绑定,远端面板没有在刷新后自动还原,新建会话可能得到新的 token 和 endpoint ID,旧任务无法自动识别它就是原来的资产连接
因此,当前方案已经完成 "任务持久化、断点协议、安全恢复",但跨刷新自动续传还需要稳定的资产身份,以及从旧 endpoint 到新会话 endpoint 的重绑定机制才能完全闭环
设计取舍与后续演进
这次改造最关键的取舍,是没有用更高并发、更复杂的持久化结构或更多 UI 状态来包装问题,而是围绕 "可信检查点" 建立最小闭环
ACK 保证进度来自目标端事实,confirmedBytes 与 checksumState 组成客户端检查点,committedBytes 代表目标端检查点,prepare 阶段的双端对账负责判断两者是否仍然对应,状态机则把恢复、失败和终态表达清楚
如果继续演进,优先级最高的不是直接放大并发数,而是先补齐稳定资产身份与 endpoint 重绑定;在此基础上,再根据真实指标决定是否引入每个 endpoint 的有限并发、源端读取复用、连接池或带宽自适应
