iOS 与传输修复:取消即时落终态 + 行操作统一 + 记录按用户分桶 + 发送端进度/完成对齐接收端
- 取消传输无效:web cancelTransfer 拆除本端 p2p/relay 后立即 completeTransfer(sessionId, "CANCELLED") 落终态,不再仅依赖服务端 transfer:state CANCELLED 回显——回显延迟、或因设备名变更致路由不中而未送达时,旧实现使记录永停“进行中”,即用户所见“取消无效”。completeTransfer 幂等,随后回显安全重入。 - 进行中/已完成传输行操作统一:进行中行新增滑动“取消”+ 长按菜单“取消/切中继”,与历史行滑动+长按“删除”对称。取消即时无确认、随 transferDone 移入 history,析构滑动动画与现实一致不抖。详情页取消/切中继入口保留。 - 传输历史/消息持久化:记录按登录用户 userId 分桶持久化(EngineController recordKey/activateRecords/persist*,含旧无桶文件一次性回退)。reset() 改为只清内存、不再 RecordsStore.clear 磁盘——会话失效(authExpired/被吊销)属凭证层事件,不应连用户数据一并销毁;同一账户重登经 activateRecords 即恢复,换账号各读各桶天然防跨账号泄露。CDropApp 在 onAppear 与 onChange(session.user.id) 激活记录。 - 发送端进度/完成对齐接收端:接收端经数据通道节流(200ms)+finalize 即时回传已收字节(ack 帧);发送端进度改用 ackedBytes(单调取大),取代仅反映本地 SCTP 缓冲抽干的 senderDeliveredBytes——避免 Windows→iOS 等慢接收端组合下发送端进度虚高领先、完成早十数秒(4MB)。streamFile 抽干 buffer 后再等 ack 追平总量才宣告完成;ackedBytes==0(旧版本接收端)走旧“抽干即完成”行为、30s 超时兜底、接收端权威 /done 经 SSE 独立完成。
This commit is contained in:
@@ -57,6 +57,14 @@ struct AppRoot: View
|
|||||||
{
|
{
|
||||||
auth.debugSession()
|
auth.debugSession()
|
||||||
}
|
}
|
||||||
|
// 冷启已有会话:立即按 userId 加载该用户的持久记录。reset() 只清内存不动磁盘,故隔夜
|
||||||
|
// 被 authExpired 清空视图后,这里读盘即恢复传输历史 / 消息(#3)。
|
||||||
|
if let uid = auth.session?.user.id { engine.activateRecords(userId: uid) }
|
||||||
|
}
|
||||||
|
// 登录完成(session 由 nil 变有)激活该用户的记录桶;authExpired / 被吊销后同一用户重登亦经此恢复。
|
||||||
|
.onChange(of: auth.session?.user.id)
|
||||||
|
{ _, uid in
|
||||||
|
if let uid { engine.activateRecords(userId: uid) }
|
||||||
}
|
}
|
||||||
// Share Extension 深链:cdrop://share 唤起 → 导入收件箱待发文件(RootView 据此弹选设备)。
|
// Share Extension 深链:cdrop://share 唤起 → 导入收件箱待发文件(RootView 据此弹选设备)。
|
||||||
// 放在 AppRoot 以便登录 / 未登录任一界面都接得到;未登录时文件留收件箱,登录后 onAppear 补扫。
|
// 放在 AppRoot 以便登录 / 未登录任一界面都接得到;未登录时文件留收件箱,登录后 onAppear 补扫。
|
||||||
|
|||||||
@@ -159,16 +159,65 @@ final class EngineController: NSObject
|
|||||||
private static let historyCap = 50
|
private static let historyCap = 50
|
||||||
private static let messagesCap = 200
|
private static let messagesCap = 200
|
||||||
|
|
||||||
|
// 当前记录归属用户:传输历史 / 消息 / 设备快照按登录用户分桶持久化(文件名带 userId,见 recordKey)。
|
||||||
|
// 这样会话失效(authExpired / 被吊销)只清内存、不毁盘,同一用户重登即恢复;换账号登录各读各的桶,
|
||||||
|
// 天然防跨账号泄露。空串=尚未激活会话(事件先于激活到达时跳过落盘,见 persist*)。
|
||||||
|
private var currentRecordUser = ""
|
||||||
|
|
||||||
override init()
|
override init()
|
||||||
{
|
{
|
||||||
super.init()
|
super.init()
|
||||||
PushRegistry.shared.engine = self
|
PushRegistry.shared.engine = self
|
||||||
// 加载上次持久化的传输历史 / 消息记录,修「重启清空」。新事件经引擎桥增量追加。
|
// 记录加载推迟到 activateRecords(会话激活时,userId 已知才知道读哪个桶)。init 阶段
|
||||||
history = RecordsStore.load([TransferItem].self, "history") ?? []
|
// auth 尚未注入、无从分桶,故此处不再加载——activateRecords 由 AppRoot 在登录 / 冷启已有
|
||||||
messages = RecordsStore.load([MessageItem].self, "messages") ?? []
|
// 会话时即时调用,加载延迟可忽略。
|
||||||
// 加载上次的设备 presence 快照:冷启 / 墓碑恢复时立即显示设备列表(而非空列表等 SSE 重连),
|
}
|
||||||
// 缓解「长时间重连」的观感;SSE 一连上即被新鲜 presence 整组替换、自校正在线状态。
|
|
||||||
devices = RecordsStore.load([DeviceItem].self, "devices") ?? []
|
// 会话激活:登录完成 / 冷启已有会话时按 userId 加载该用户的持久记录到内存。reset() 只清内存不动
|
||||||
|
// 磁盘,故隔夜被 authExpired 清空视图后,同一用户重登经此重新读盘即恢复历史 / 消息。幂等:重复激活
|
||||||
|
// 仅重读(内存与磁盘一致,无副作用)。legacy 回退:旧版本无分桶的 "history"/"messages"/"devices"
|
||||||
|
// 文件首次激活时一并采纳(随后落盘即写入分桶键),避免升级丢失既有记录。
|
||||||
|
func activateRecords(userId: String)
|
||||||
|
{
|
||||||
|
guard !userId.isEmpty else { return }
|
||||||
|
currentRecordUser = userId
|
||||||
|
history = RecordsStore.load([TransferItem].self, recordKey("history"))
|
||||||
|
?? RecordsStore.load([TransferItem].self, "history") ?? []
|
||||||
|
messages = RecordsStore.load([MessageItem].self, recordKey("messages"))
|
||||||
|
?? RecordsStore.load([MessageItem].self, "messages") ?? []
|
||||||
|
devices = RecordsStore.load([DeviceItem].self, recordKey("devices"))
|
||||||
|
?? RecordsStore.load([DeviceItem].self, "devices") ?? []
|
||||||
|
}
|
||||||
|
|
||||||
|
// 分桶文件键:userId(broker subject,通常即 UUID 式安全字符)拼到记录名前。保险起见把非
|
||||||
|
// 字母数字 / - / _ 的字符替换为 _,确保落地为合法文件名。
|
||||||
|
private func recordKey(_ name: String) -> String
|
||||||
|
{
|
||||||
|
let safe = String(currentRecordUser.map
|
||||||
|
{ ch in
|
||||||
|
(ch.isLetter || ch.isNumber || ch == "-" || ch == "_") ? ch : "_"
|
||||||
|
})
|
||||||
|
return "\(safe)-\(name)"
|
||||||
|
}
|
||||||
|
|
||||||
|
// 落盘助手:仅在已激活会话(currentRecordUser 非空)时写当前用户的桶。未激活时跳过,避免把
|
||||||
|
// 增量事件写进 "-name" 的无主桶。
|
||||||
|
private func persistHistory()
|
||||||
|
{
|
||||||
|
guard !currentRecordUser.isEmpty else { return }
|
||||||
|
RecordsStore.save(history, recordKey("history"))
|
||||||
|
}
|
||||||
|
|
||||||
|
private func persistMessages()
|
||||||
|
{
|
||||||
|
guard !currentRecordUser.isEmpty else { return }
|
||||||
|
RecordsStore.save(messages, recordKey("messages"))
|
||||||
|
}
|
||||||
|
|
||||||
|
private func persistDevices()
|
||||||
|
{
|
||||||
|
guard !currentRecordUser.isEmpty else { return }
|
||||||
|
RecordsStore.save(devices, recordKey("devices"))
|
||||||
}
|
}
|
||||||
|
|
||||||
// makeWebView:构建离屏 WebView——注入 __CDROP_BOOT__(device_type:"ios")、注册消息
|
// makeWebView:构建离屏 WebView——注入 __CDROP_BOOT__(device_type:"ios")、注册消息
|
||||||
@@ -232,10 +281,11 @@ final class EngineController: NSObject
|
|||||||
seenIncomingTransfers = []
|
seenIncomingTransfers = []
|
||||||
activeTab = "transfer"
|
activeTab = "transfer"
|
||||||
pendingLoginRequest = nil
|
pendingLoginRequest = nil
|
||||||
// 登出清空持久记录,不跨账号残留(重启保活仅针对同一登录态)。
|
// 不再清磁盘记录:传输历史 / 消息是用户数据,会话失效(authExpired / 被吊销)属凭证层事件,
|
||||||
RecordsStore.clear("history")
|
// 不应连用户数据一起销毁——否则隔夜令牌过期一次就抹掉全部历史(#3 根因)。这里只断开会话归属
|
||||||
RecordsStore.clear("messages")
|
// (置空 currentRecordUser),磁盘按 userId 分桶留存:同一用户重登经 activateRecords 即恢复,
|
||||||
RecordsStore.clear("devices")
|
// 换账号登录各读各的桶不会串。要彻底清除走设置页登出后用记录页 / 消息页的「清空」(用户主动)。
|
||||||
|
currentRecordUser = ""
|
||||||
status = t("ios.engine.disconnected")
|
status = t("ios.engine.disconnected")
|
||||||
deviceName = ""
|
deviceName = ""
|
||||||
PushRegistry.shared.reset()
|
PushRegistry.shared.reset()
|
||||||
@@ -251,26 +301,26 @@ final class EngineController: NSObject
|
|||||||
func deleteTransferRecord(_ sessionId: String)
|
func deleteTransferRecord(_ sessionId: String)
|
||||||
{
|
{
|
||||||
history.removeAll { $0.sessionId == sessionId }
|
history.removeAll { $0.sessionId == sessionId }
|
||||||
RecordsStore.save(history, "history")
|
persistHistory()
|
||||||
}
|
}
|
||||||
|
|
||||||
func clearHistory()
|
func clearHistory()
|
||||||
{
|
{
|
||||||
history.removeAll()
|
history.removeAll()
|
||||||
RecordsStore.clear("history")
|
RecordsStore.clear(recordKey("history"))
|
||||||
}
|
}
|
||||||
|
|
||||||
// 删除一条消息(H)/ 清空全部(持久化同步)。
|
// 删除一条消息(H)/ 清空全部(持久化同步)。
|
||||||
func deleteMessage(_ id: String)
|
func deleteMessage(_ id: String)
|
||||||
{
|
{
|
||||||
messages.removeAll { $0.id == id }
|
messages.removeAll { $0.id == id }
|
||||||
RecordsStore.save(messages, "messages")
|
persistMessages()
|
||||||
}
|
}
|
||||||
|
|
||||||
func clearMessages()
|
func clearMessages()
|
||||||
{
|
{
|
||||||
messages.removeAll()
|
messages.removeAll()
|
||||||
RecordsStore.clear("messages")
|
RecordsStore.clear(recordKey("messages"))
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---- 瞬时通知(#2) ----
|
// ---- 瞬时通知(#2) ----
|
||||||
@@ -655,7 +705,7 @@ extension EngineController: WKScriptMessageHandler
|
|||||||
if let p = payload as? [String: Any], let raw = p["devices"] as? [Any]
|
if let p = payload as? [String: Any], let raw = p["devices"] as? [Any]
|
||||||
{
|
{
|
||||||
devices = raw.compactMap { ($0 as? [String: Any]).flatMap(Self.parseDevice) }
|
devices = raw.compactMap { ($0 as? [String: Any]).flatMap(Self.parseDevice) }
|
||||||
RecordsStore.save(devices, "devices") // 持久化 presence 快照,供下次冷启即时显示
|
persistDevices() // 持久化 presence 快照(按用户分桶),供下次冷启即时显示
|
||||||
}
|
}
|
||||||
case "transfers":
|
case "transfers":
|
||||||
if let p = payload as? [String: Any], let raw = p["active"] as? [Any]
|
if let p = payload as? [String: Any], let raw = p["active"] as? [Any]
|
||||||
@@ -678,7 +728,7 @@ extension EngineController: WKScriptMessageHandler
|
|||||||
{
|
{
|
||||||
messages.removeLast(messages.count - Self.messagesCap)
|
messages.removeLast(messages.count - Self.messagesCap)
|
||||||
}
|
}
|
||||||
RecordsStore.save(messages, "messages")
|
persistMessages()
|
||||||
// 新到的对端消息:不在消息页时计未读红点(#5)。本机发出的乐观回显(outgoing)不计。
|
// 新到的对端消息:不在消息页时计未读红点(#5)。本机发出的乐观回显(outgoing)不计。
|
||||||
if isNew, item.direction == "incoming", activeTab != "messages" { unreadMessages += 1 }
|
if isNew, item.direction == "incoming", activeTab != "messages" { unreadMessages += 1 }
|
||||||
}
|
}
|
||||||
@@ -692,7 +742,7 @@ extension EngineController: WKScriptMessageHandler
|
|||||||
history.insert(item, at: 0)
|
history.insert(item, at: 0)
|
||||||
if history.count > Self.historyCap { history.removeLast(history.count - Self.historyCap) }
|
if history.count > Self.historyCap { history.removeLast(history.count - Self.historyCap) }
|
||||||
syncBackgroundTask()
|
syncBackgroundTask()
|
||||||
RecordsStore.save(history, "history") // 持久化,跨重启存活
|
persistHistory() // 持久化(按用户分桶),跨重启存活
|
||||||
}
|
}
|
||||||
case "sendStarted":
|
case "sendStarted":
|
||||||
// 随后的 transfers 事件会带出该活跃项,这里无需额外处理。
|
// 随后的 transfers 事件会带出该活跃项,这里无需额外处理。
|
||||||
|
|||||||
@@ -227,9 +227,10 @@ struct TransferListView: View
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 传输行:卡片观感不变;点按进详情;历史记录支持滑动删除 + 长按删除(与设备 / 文件长条一致)。
|
// 传输行:卡片观感不变;点按进详情;所有长条行操作逻辑统一——滑动 + 长按都给主操作(与设备 /
|
||||||
// 进行中的传输不挂删除手势(取消 / 切中继在详情页)。删除是真删除(无需二次确认),故析构滑动按钮
|
// 文件长条一致)。历史记录=删除(真删除、无需二次确认);进行中=取消(即时、无确认),长按菜单
|
||||||
// 的删除动画与实际一致、不会抖动。
|
// 另给「切中继」。两类的滑动 / 长按都直挂破坏性操作,析构滑动动画与实际一致(取消即随 transferDone
|
||||||
|
// 移入 history、删除即移除),不会抖动。详情页仍保留取消 / 切中继入口。
|
||||||
@ViewBuilder
|
@ViewBuilder
|
||||||
private func transferRow(_ item: TransferItem, history: Bool) -> some View
|
private func transferRow(_ item: TransferItem, history: Bool) -> some View
|
||||||
{
|
{
|
||||||
@@ -256,6 +257,22 @@ struct TransferListView: View
|
|||||||
else
|
else
|
||||||
{
|
{
|
||||||
card
|
card
|
||||||
|
.swipeActions(edge: .trailing)
|
||||||
|
{
|
||||||
|
Button(role: .destructive) { engine.cancelTransfer(item.sessionId) }
|
||||||
|
label: { Label(t("ios.transfer.cancel"), systemImage: "xmark.circle") }
|
||||||
|
}
|
||||||
|
.contextMenu
|
||||||
|
{
|
||||||
|
Button(role: .destructive) { engine.cancelTransfer(item.sessionId) }
|
||||||
|
label: { Label(t("ios.transfer.cancel"), systemImage: "xmark.circle") }
|
||||||
|
if item.mode != "relay"
|
||||||
|
{
|
||||||
|
Button { engine.switchToRelay(item.sessionId) }
|
||||||
|
label: { Label(t("ios.transfer.forceRelay"),
|
||||||
|
systemImage: "antenna.radiowaves.left.and.right") }
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -30,6 +30,12 @@ const LOW_WATERMARK = 4 * 1024 * 1024;
|
|||||||
// 大小与流控(HIGH/LOW watermark)完全不变,接收端无感。
|
// 大小与流控(HIGH/LOW watermark)完全不变,接收端无感。
|
||||||
const READ_BLOCK_SIZE = 4 * 1024 * 1024;
|
const READ_BLOCK_SIZE = 4 * 1024 * 1024;
|
||||||
const ICE_CANDIDATE_POOL = 4;
|
const ICE_CANDIDATE_POOL = 4;
|
||||||
|
// 接收端把"已实际收下字节数"回传发送端的节流间隔:每 chunk 都回会与数据争用 SCTP、拖慢吞吐,
|
||||||
|
// 故尾沿合并到 ~200ms 发一次最新值。
|
||||||
|
const ACK_INTERVAL_MS = 200;
|
||||||
|
// 发送端等接收端 ack 追平总量后才宣告完成的上限:超时仍未追平则照常收尾(接收端权威 /done 经
|
||||||
|
// SSE 也会独立完成),避免 ack 丢失 / 接收端异常致发送端悬挂在"完成中"。
|
||||||
|
const ACK_COMPLETE_TIMEOUT_MS = 30_000;
|
||||||
|
|
||||||
export interface FileMeta
|
export interface FileMeta
|
||||||
{
|
{
|
||||||
@@ -100,6 +106,12 @@ class Session
|
|||||||
// 发送进度轮询:streamFile 期间每 250ms 刷一次,让"已交付字节"在 backpressure
|
// 发送进度轮询:streamFile 期间每 250ms 刷一次,让"已交付字节"在 backpressure
|
||||||
// 等待期间也能平滑爬升、不误报 stalled(见 streamFile / senderDeliveredBytes)。
|
// 等待期间也能平滑爬升、不误报 stalled(见 streamFile / senderDeliveredBytes)。
|
||||||
private sendProgressPollId: number | null = null;
|
private sendProgressPollId: number | null = null;
|
||||||
|
// 接收端回传的"已实际收下字节数"(ack):发送端进度 / 完成判定的权威依据,取代仅反映本地 SCTP
|
||||||
|
// 缓冲抽干的 senderDeliveredBytes——避免接收端(WK/iOS)慢速落盘时发送端进度虚高、抢先完成。
|
||||||
|
private ackedBytes = 0;
|
||||||
|
// 接收端 ack 回传节流(尾沿合并,见 scheduleAck)。
|
||||||
|
private lastAckSentAt = 0;
|
||||||
|
private ackTimer: number | null = null;
|
||||||
|
|
||||||
constructor(public sessionId: string, public peerName: string, public role: "sender" | "receiver")
|
constructor(public sessionId: string, public peerName: string, public role: "sender" | "receiver")
|
||||||
{
|
{
|
||||||
@@ -246,11 +258,23 @@ class Session
|
|||||||
private emitProgress(): void
|
private emitProgress(): void
|
||||||
{
|
{
|
||||||
const total = this.role === "sender" ? this.outgoingTotal : (this.receivedMeta?.size ?? 0);
|
const total = this.role === "sender" ? this.outgoingTotal : (this.receivedMeta?.size ?? 0);
|
||||||
const bytes = this.role === "sender" ? this.senderDeliveredBytes() : this.receivedBytes;
|
const bytes = this.role === "sender" ? this.senderProgressBytes() : this.receivedBytes;
|
||||||
for (const cb of this.progressListeners) { cb({ bytes, total }); }
|
for (const cb of this.progressListeners) { cb({ bytes, total }); }
|
||||||
this.throttledPushUpdate();
|
this.throttledPushUpdate();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 发送端进度字节:优先用接收端回传的已收字节(ackedBytes)——它反映接收端真实收程,避免
|
||||||
|
* bufferedAmount 仅代表本地 SCTP 缓冲抽干、而接收端(如 WK/iOS)仍在慢速落盘时进度虚高领先
|
||||||
|
* (4MB 文件抢先十数秒的失真根因)。接收端首个 ack 到达前回退到本地交付估计,使连接初期仍有
|
||||||
|
* 进度显示。封顶 outgoingTotal,防 ack 异常越界。
|
||||||
|
*/
|
||||||
|
private senderProgressBytes(): number
|
||||||
|
{
|
||||||
|
if (this.ackedBytes > 0) { return Math.min(this.outgoingTotal, this.ackedBytes); }
|
||||||
|
return this.senderDeliveredBytes();
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 发送端"已交付字节"=已离开本地 buffer 的量(bytesSent − bufferedAmount)。
|
* 发送端"已交付字节"=已离开本地 buffer 的量(bytesSent − bufferedAmount)。
|
||||||
* bytesSent 只是"已塞进 dc 的本地 16MiB buffer",会被 64KB 一片秒满到 16MiB 后
|
* bytesSent 只是"已塞进 dc 的本地 16MiB buffer",会被 64KB 一片秒满到 16MiB 后
|
||||||
@@ -284,6 +308,22 @@ class Session
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 等接收端 ack 追平总量(确认已收全)再让发送端完成。超时兜底:ACK_COMPLETE_TIMEOUT_MS 内仍
|
||||||
|
// 未追平则返回照常收尾——本地 buffer 已抽干(字节确已离开本端)、接收端权威 /done 经 SSE 也会
|
||||||
|
// 独立驱动 store 终态,避免 ack 丢失 / 接收端异常致发送端悬挂。取消时立即返回。
|
||||||
|
private async waitForReceiverAck(total: number): Promise<void>
|
||||||
|
{
|
||||||
|
// buffer 抽干后仍无任何 ack=接收端不支持回传(旧版本):保持旧行为"抽干即完成",不空等
|
||||||
|
// 超时。新接收端在抽干前必已回过若干 ack(200ms 节流 + 抽干本身要等 SCTP 送达慢接收端),
|
||||||
|
// 故 ackedBytes>0 即可靠区分新旧。极小文件抽干快于首个 ack 时同样走此分支,无碍。
|
||||||
|
if (this.ackedBytes === 0) { return; }
|
||||||
|
const deadline = Date.now() + ACK_COMPLETE_TIMEOUT_MS;
|
||||||
|
while (this.ackedBytes < total && !this.canceled && Date.now() < deadline)
|
||||||
|
{
|
||||||
|
await new Promise((r) => window.setTimeout(r, 100));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private classifyAndCount(c: RTCIceCandidate | RTCIceCandidateInit, into: CandidateBreakdown): void
|
private classifyAndCount(c: RTCIceCandidate | RTCIceCandidateInit, into: CandidateBreakdown): void
|
||||||
{
|
{
|
||||||
const sdp = (c as RTCIceCandidate).candidate ?? (c as RTCIceCandidateInit).candidate ?? "";
|
const sdp = (c as RTCIceCandidate).candidate ?? (c as RTCIceCandidateInit).candidate ?? "";
|
||||||
@@ -349,6 +389,16 @@ class Session
|
|||||||
this.setPhase("completing");
|
this.setPhase("completing");
|
||||||
void this.finalizeIncoming();
|
void this.finalizeIncoming();
|
||||||
}
|
}
|
||||||
|
else if (msg.type === "ack")
|
||||||
|
{
|
||||||
|
// 接收端回传的"已收下字节数":仅发送端处理。单调取大(防乱序帧把进度回退),
|
||||||
|
// 据此把发送端进度 / 完成判定对齐接收端真实收程。
|
||||||
|
if (this.role === "sender" && typeof msg.bytes === "number")
|
||||||
|
{
|
||||||
|
this.ackedBytes = Math.max(this.ackedBytes, msg.bytes);
|
||||||
|
this.emitProgress();
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
catch
|
catch
|
||||||
{
|
{
|
||||||
@@ -368,9 +418,35 @@ class Session
|
|||||||
return sink;
|
return sink;
|
||||||
});
|
});
|
||||||
this.emitProgress();
|
this.emitProgress();
|
||||||
|
// 节流回传已收字节给发送端,使其进度 / 完成判定对齐接收端真实收程(见 scheduleAck)。
|
||||||
|
this.scheduleAck();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 接收端把已收下的累计字节数(receivedBytes)经数据通道回传发送端。节流到 ACK_INTERVAL_MS、
|
||||||
|
// 尾沿合并发最新值——每 chunk 都回会与下行数据争用 SCTP 通道、反噬吞吐。仅 receiver 起作用。
|
||||||
|
private scheduleAck(): void
|
||||||
|
{
|
||||||
|
if (this.role !== "receiver" || this.ackTimer !== null) { return; }
|
||||||
|
const wait = Math.max(0, ACK_INTERVAL_MS - (Date.now() - this.lastAckSentAt));
|
||||||
|
this.ackTimer = window.setTimeout(() =>
|
||||||
|
{
|
||||||
|
this.ackTimer = null;
|
||||||
|
this.sendAck();
|
||||||
|
}, wait);
|
||||||
|
}
|
||||||
|
|
||||||
|
// 立即回传一次当前 receivedBytes(finalize 时调用,确保发送端及时见到 ackedBytes==total)。
|
||||||
|
private sendAck(): void
|
||||||
|
{
|
||||||
|
if (this.role !== "receiver") { return; }
|
||||||
|
const dc = this.dc;
|
||||||
|
if (!dc || dc.readyState !== "open") { return; }
|
||||||
|
this.lastAckSentAt = Date.now();
|
||||||
|
try { dc.send(JSON.stringify({ type: "ack", bytes: this.receivedBytes })); }
|
||||||
|
catch { /* 通道抖动:下个 chunk 的 scheduleAck 会补发最新值 */ }
|
||||||
|
}
|
||||||
|
|
||||||
/** 收方看门狗:超过 5s 没收到任何消息(onmessage 没 fire)即 console.warn,
|
/** 收方看门狗:超过 5s 没收到任何消息(onmessage 没 fire)即 console.warn,
|
||||||
* 方便诊断主线程被微任务堵住的情形。仅 receiver 角色起作用。 */
|
* 方便诊断主线程被微任务堵住的情形。仅 receiver 角色起作用。 */
|
||||||
private startReceiverWatchdog(): void
|
private startReceiverWatchdog(): void
|
||||||
@@ -403,6 +479,9 @@ class Session
|
|||||||
private async finalizeIncoming(): Promise<void>
|
private async finalizeIncoming(): Promise<void>
|
||||||
{
|
{
|
||||||
if (!this.receivedMeta) { return; }
|
if (!this.receivedMeta) { return; }
|
||||||
|
// 收到 done=所有 chunk 已到(dc ordered),receivedBytes 已等于 total。立即回传一次最终
|
||||||
|
// ack,使发送端及时见到 ackedBytes==total 而完成,不必等本端 sink.close 落盘或下次节流。
|
||||||
|
this.sendAck();
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
const sink = await this.receiveQueue;
|
const sink = await this.receiveQueue;
|
||||||
@@ -568,12 +647,13 @@ class Session
|
|||||||
|
|
||||||
dc.send(JSON.stringify({ type: "done" }));
|
dc.send(JSON.stringify({ type: "done" }));
|
||||||
this.setPhase("completing");
|
this.setPhase("completing");
|
||||||
// 抽干整个 buffer(数据 + done 帧)再宣告完成:dc.send 只是入本地 buffer,真正
|
// 抽干本地 buffer(数据 + done 帧)只代表字节离开本端 SCTP,不代表接收端已收下——接收端
|
||||||
// 送达要等 SCTP 发出。早先送完入 buffer 即 setState("completed"),会在还有多达
|
// (如 WK/iOS)可能仍在慢速 onmessage / 落盘。早先"buffer 抽干即 setState completed"会在
|
||||||
// 16MiB 在途时就标记完成,与接收端"仍在下载"不一致。等 bufferedAmount 落到 0=
|
// 接收端尚有大量未处理时就抢先完成(4MB 文件可领先十数秒、进度虚高)。故先抽干,再等接收端
|
||||||
// 字节已全部交付网络,接收端随后 finalize 并 POST /done(那条权威 /done 仍由接收
|
// 回传的 ack 追平总量=其确已收全,发送端才宣告完成,与接收端同步。进度轮询在等待期持续刷新
|
||||||
// 端发,这里只在抽干后驱动发送端自身 UI 收尾)。
|
// (ackedBytes 爬升可见);超时兜底见 waitForReceiverAck。权威 /done 仍由接收端 finalize 后发。
|
||||||
await waitForBufferLow(dc, 0);
|
await waitForBufferLow(dc, 0);
|
||||||
|
await this.waitForReceiverAck(src.size);
|
||||||
this.stopSendProgressPoll();
|
this.stopSendProgressPoll();
|
||||||
this.bytesSent = src.size;
|
this.bytesSent = src.size;
|
||||||
this.emitProgress();
|
this.emitProgress();
|
||||||
@@ -733,6 +813,11 @@ class Session
|
|||||||
this.selectedPairPollId = null;
|
this.selectedPairPollId = null;
|
||||||
}
|
}
|
||||||
this.stopSendProgressPoll();
|
this.stopSendProgressPoll();
|
||||||
|
if (this.ackTimer !== null)
|
||||||
|
{
|
||||||
|
window.clearTimeout(this.ackTimer);
|
||||||
|
this.ackTimer = null;
|
||||||
|
}
|
||||||
this.dc?.close();
|
this.dc?.close();
|
||||||
this.pc.close();
|
this.pc.close();
|
||||||
this.setState("closed");
|
this.setState("closed");
|
||||||
|
|||||||
@@ -356,6 +356,12 @@ export async function cancelTransfer(sessionId: string): Promise<void>
|
|||||||
if (inc) { p2pCleanup(inc.sender); }
|
if (inc) { p2pCleanup(inc.sender); }
|
||||||
cleanupTransfer(sessionId);
|
cleanupTransfer(sessionId);
|
||||||
|
|
||||||
|
// 本端立即落终态:把记录移入 history(取消即时反映到 UI),不依赖服务端 transfer:state
|
||||||
|
// CANCELLED 回显。旧实现仅靠该回显才 completeTransfer,回显延迟 / 因设备名变更等未送达时,
|
||||||
|
// 记录会永远停在「进行中」——即用户所见「取消无效」。本端 p2p / relay 已就地拆除,记录确已
|
||||||
|
// 终结,故先行落终态正确;completeTransfer 幂等(已移走即 no-op),随后回显安全重入。
|
||||||
|
useAppStore.getState().completeTransfer(sessionId, "CANCELLED");
|
||||||
|
|
||||||
const r = await apiFetch(`/api/transfer/${sessionId}/cancel`, { method: "POST" });
|
const r = await apiFetch(`/api/transfer/${sessionId}/cancel`, { method: "POST" });
|
||||||
if (!r.ok && r.status !== 409)
|
if (!r.ok && r.status !== 409)
|
||||||
{
|
{
|
||||||
|
|||||||
Reference in New Issue
Block a user