import Foundation import WebRTC // LibWebRtcEngine —— iOS 原生数据面引擎,用 libwebrtc(Safari 同款 RTCPeerConnection)在原生进程内 // 完成 P2P 文件收发。取代早先 gomobile/pion 路径:pion 的 Go raw socket 不与 iOS Network framework // 集成(WebKit/libwebrtc 才集成),真机确证架构性不可靠(见 desktop/NATIVE-TRANSFER.md §7)。 // // 线协议与 desktop/engine(pion)及 web/src/features/transfer/p2p.ts 逐字节一致,故可与桌面 pion // 引擎、浏览器 / 旧 iOS 的 JS 引擎互通:DataChannel "commilitia-drop-file"(ordered),控制帧 // 文本帧、文件分片走二进制帧(64KB),信令 payload 形如 {type, sdp?, candidate?}。本类是 session.go + // engine.go 的忠实 Swift 端口——逻辑、水位、ack 语义、收尾顺序一一对应。 // // 分工同桌面:本引擎只跑纯 WebRTC + 文件 I/O,一切 HTTP(出站信令 POST / 状态机 POST)留 JS // (见 web/src/features/transfer/p2pIos.ts)。事件经 P2PEngineDelegate 回灌——宿主即 EngineController // 经 cdropEngine 桥(P2PEventsBridge)转给无头 JS 引擎。 // // 宿主契约(EngineController 的 4 个 p2p RPC + net/ios.ts 的 nativeP2P*)一字不改;只把 Swift 侧引擎 // 从 gomobile EngineEngine 换成本类。翻 web 侧 IOS_NATIVE=true 即启用。 // MARK: - 事件回调(宿主实现,等价 desktop/engine 的 Events 接口) protocol P2PEngineDelegate: AnyObject { // 已传字节:发送端=接收端 ack 追平值(无 ack 时回退本地已交付估计),接收端=已写盘字节。 func engineDidProgress(sessionId: String, bytes: Int64) // 会话状态:connected / completed / failed / closed。 func engineDidChangeState(sessionId: String, state: String) // 出站信令(offer/answer/ice):payloadJSON 为信令 JSON 串,宿主 POST /api/hub/signal。 func engineDidEmitSignal(sessionId: String, toPeer: String, payloadJSON: String) // 接收端落盘完成的绝对路径(沙盒 Documents)。 func engineDidSave(sessionId: String, path: String) // 选中候选对(形如 "host/udp"),供调试面板展示实际连通路径。 func engineDidSelectIcePair(sessionId: String, local: String, remote: String) // 诊断日志(候选 / ICE 态 / 选中对 / 发送停滞),形如 "session ...",供日志界面按会话聚合。 func engineDidLog(_ line: String) } // MARK: - 线协议常量(与 desktop/engine wire.go、web p2p.ts 一致) private enum Wire { static let channelName = "commilitia-drop-file" static let chunkSize = 64 * 1024 static let highWatermark: UInt64 = 16 * 1024 * 1024 static let lowWatermark: UInt64 = 4 * 1024 * 1024 static let ackIntervalMs = 200 static let progressThrottleNs: Int64 = 100 * 1_000_000 // ~10Hz,护宿主主线程(每片过桥 evaluateJavaScript) static let ackCompleteTimeout: TimeInterval = 30 // ICE 候选预热池,与 web p2p.ts ICE_CANDIDATE_POOL 对齐。 static let iceCandidatePool: Int32 = 4 } struct EngineError: LocalizedError { let msg: String init(_ m: String) { msg = m } var errorDescription: String? { msg } } // MARK: - 引擎(持有所有活跃会话,按对端设备名索引,等价 engine.go 的 Engine) final class LibWebRtcEngine { // 单次 SSL 初始化(libwebrtc 要求)。全局 let 惰性、线程安全地执行一次。 private static let sslReady: Bool = RTCInitializeSSL() private let factory: RTCPeerConnectionFactory // 内部可设:生产由 init 注入(EngineController 的 P2PEventsBridge);测试可重定向到环回宿主。 weak var delegate: P2PEngineDelegate? private let lock = NSLock() private var byPeer: [String: P2PTransferSession] = [:] private var downloadDirValue: String init(downloadDir: String, delegate: P2PEngineDelegate) { _ = LibWebRtcEngine.sslReady factory = RTCPeerConnectionFactory() downloadDirValue = downloadDir self.delegate = delegate } var events: P2PEngineDelegate? { delegate } func downloadDir() -> String { lock.lock(); defer { lock.unlock() } return downloadDirValue } func setDownloadDir(_ dir: String) { lock.lock(); downloadDirValue = dir; lock.unlock() } // StartOutgoing:建 PeerConnection + DataChannel,发 offer,通道就绪后流式发出 filePath。 // filePath 为本机绝对路径(原生从 commilitia-drop-file:// 解析得到)。completion 在 offer 已铸并 // setLocalDescription 完成后回(或出错时回错误),对齐 JS 侧信令链 seed 语义(见 p2pIos.ts chain)。 func startOutgoing(sessionId: String, peerName: String, filePath: String, iceServersJSON: String, completion: @escaping (Error?) -> Void) { let fm = FileManager.default var isDir: ObjCBool = false guard fm.fileExists(atPath: filePath, isDirectory: &isDir), !isDir.boolValue else { completion(EngineError("not a file: \(filePath)")); return } let size = ((try? fm.attributesOfItem(atPath: filePath))?[.size] as? NSNumber)?.int64Value ?? 0 let s = P2PTransferSession(engine: self, factory: factory, sessionId: sessionId, peerName: peerName, role: .sender) s.filePath = filePath s.fileSize = size do { try s.createPeerConnection(iceServersJSON) } catch { completion(error); return } put(peerName, s) do { try s.createSenderDataChannel() } catch { s.fail("create datachannel: \(error.localizedDescription)"); completion(error); return } s.createOfferAndEmit(completion: completion) } // StartIncoming:建 PeerConnection 等对端 offer,收到 DataChannel 后按线协议落盘。须在对端 // offer 到达前调用。completion 立即回(会话已登记,后续入站 offer/ice 可路由)。 func startIncoming(sessionId: String, peerName: String, iceServersJSON: String, completion: @escaping (Error?) -> Void) { let s = P2PTransferSession(engine: self, factory: factory, sessionId: sessionId, peerName: peerName, role: .receiver) do { try s.createPeerConnection(iceServersJSON) } catch { completion(error); return } put(peerName, s) completion(nil) } // HandleSignal:把宿主从 SSE 收到的入站信令喂给对应会话(按对端设备名路由)。completion 在信令 // 应用完成后回——对 offer 是 answer 已铸并发出后,对 answer 是 setRemoteDescription 完成后,对 // ice 是 addIceCandidate 完成后。这保证 JS 侧信令链按序投递(offer/answer 先于其后 trickle 候选)。 func handleSignal(fromPeer: String, payloadJSON: String, completion: @escaping (Error?) -> Void) { lock.lock(); let s = byPeer[fromPeer]; lock.unlock() guard let s else { completion(EngineError("no active session for peer \(fromPeer)")); return } s.handleSignal(payloadJSON, completion: completion) } // Cancel:拆除指定会话(取消 / 中继回退时由宿主调用)。静默拆除,不回 OnState。 func cancel(sessionId: String) { lock.lock() var target: P2PTransferSession? for (_, s) in byPeer where s.sessionId == sessionId { target = s; break } lock.unlock() target?.cleanup() } fileprivate func put(_ peer: String, _ s: P2PTransferSession) { lock.lock() let old = byPeer[peer] byPeer[peer] = s lock.unlock() // 同一对端若已有会话,先拆旧的(沿用 web / 桌面按对端覆盖的语义)。异步拆,避免在 RPC / 主线程 // 上同步 close(对齐 engine.go 的 `go old.cleanup()`)。 if let old, old !== s { DispatchQueue.global(qos: .utility).async { old.cleanup() } } } fileprivate func remove(_ peer: String, _ s: P2PTransferSession) { lock.lock() if byPeer[peer] === s { byPeer.removeValue(forKey: peer) } lock.unlock() } } // MARK: - 单条传输会话(pc + dc + 角色相关收发循环,等价 session.go 的 session) final class P2PTransferSession: NSObject { enum Role { case sender, receiver } let sessionId: String let peerName: String let role: Role private weak var engine: LibWebRtcEngine? private let factory: RTCPeerConnectionFactory private var pc: RTCPeerConnection? private var dc: RTCDataChannel? // 发送端 var filePath = "" var fileSize: Int64 = 0 // 计数器(跨线程:发送循环 / pion 回调 / dc 委托)。cLock 守护。 private let cLock = NSLock() private var bytesSent: Int64 = 0 // 已塞进本地 SCTP 缓冲的字节 private var receivedBytes: Int64 = 0 // 接收端已写盘字节 private var ackedBytes: Int64 = 0 // 接收端经 ack 回传的已收字节 private var lastEmitNs: Int64 = 0 // 进度节流时间戳 // 接收端落盘。outLock 守护 outHandle 的读/写/关——writeChunk 在 libwebrtc 信令线程跑,cleanup 可能 // 在主线程(JS p2pCancel / 中继回退 / SSE 终态)并发关句柄;无锁则 ARC retain/release 与 close/write // 同体竞争 → 可能崩溃(Go 的 *os.File 有内建 fdMutex 故安全,Swift FileHandle 无此保证)。 private var outHandle: FileHandle? private var outName = "" private var tmpPath = "" private let outLock = NSLock() // 发送流 / 背压 private let streamQueue: DispatchQueue private let bufSem = DispatchSemaphore(value: 0) // didChangeBufferedAmount / ack 唤醒等待循环 private var ackTimer: DispatchSourceTimer? private var senderStarted = false // 生命周期 private let stateLock = NSLock() private var state = "" private var doneFlag = false private var cleanedUp = false init(engine: LibWebRtcEngine, factory: RTCPeerConnectionFactory, sessionId: String, peerName: String, role: Role) { self.engine = engine self.factory = factory self.sessionId = sessionId self.peerName = peerName self.role = role streamQueue = DispatchQueue(label: "net.commilitia.cdrop.p2p.\(sessionId)") super.init() } // MARK: 建连 func createPeerConnection(_ iceServersJSON: String) throws { let config = RTCConfiguration() config.iceServers = Self.parseICEServers(iceServersJSON) config.sdpSemantics = .unifiedPlan config.continualGatheringPolicy = .gatherOnce config.iceCandidatePoolSize = Wire.iceCandidatePool // 默认网络处理(不禁链路本地、不强 IPv4、不限接口):libwebrtc 与 iOS Network framework 原生 // 集成,同内网经 host↔host 直连。绝不在此加 iOS 定向过滤——曾把 IPv4-only 漏进桌面共享引擎致 // 桌面↔iOS-WebKit 的 IPv6 host 对被裁、退中继(重大回归)。 let constraints = RTCMediaConstraints(mandatoryConstraints: nil, optionalConstraints: nil) guard let pc = factory.peerConnection(with: config, constraints: constraints, delegate: self) else { throw EngineError("failed to create peer connection") } self.pc = pc } func createSenderDataChannel() throws { guard let pc else { throw EngineError("no pc") } let dcConfig = RTCDataChannelConfiguration() dcConfig.isOrdered = true guard let dc = pc.dataChannel(forLabel: Wire.channelName, configuration: dcConfig) else { throw EngineError("failed to create data channel") } dc.delegate = self self.dc = dc } func createOfferAndEmit(completion: @escaping (Error?) -> Void) { guard let pc else { completion(EngineError("no pc")); return } let constraints = RTCMediaConstraints(mandatoryConstraints: nil, optionalConstraints: nil) pc.offer(for: constraints) { [weak self] sdp, err in guard let self else { completion(EngineError("session gone")); return } if let err { self.fail("create offer: \(err.localizedDescription)"); completion(err); return } guard let sdp else { self.fail("create offer: nil"); completion(EngineError("nil offer")); return } pc.setLocalDescription(sdp) { [weak self] err2 in guard let self else { completion(EngineError("session gone")); return } if let err2 { self.fail("set local (offer): \(err2.localizedDescription)"); completion(err2); return } self.emitSignal(sdp: sdp, candidate: nil) completion(nil) } } } // MARK: 入站信令 func handleSignal(_ payloadJSON: String, completion: @escaping (Error?) -> Void) { guard let pc else { completion(EngineError("no pc")); return } guard let data = payloadJSON.data(using: .utf8), let obj = (try? JSONSerialization.jsonObject(with: data)) as? [String: Any], let type = obj["type"] as? String else { completion(EngineError("decode signal")); return } log("signal in: \(type)") switch type { case "offer": guard let sdpObj = obj["sdp"] as? [String: Any], let sdpStr = sdpObj["sdp"] as? String else { completion(EngineError("offer missing sdp")); return } let desc = RTCSessionDescription(type: .offer, sdp: sdpStr) pc.setRemoteDescription(desc) { [weak self] err in guard let self else { completion(EngineError("gone")); return } if let err { completion(err); return } let constraints = RTCMediaConstraints(mandatoryConstraints: nil, optionalConstraints: nil) pc.answer(for: constraints) { [weak self] ans, aerr in guard let self else { completion(EngineError("gone")); return } if let aerr { completion(aerr); return } guard let ans else { completion(EngineError("nil answer")); return } pc.setLocalDescription(ans) { [weak self] lerr in guard let self else { completion(EngineError("gone")); return } if let lerr { completion(lerr); return } self.emitSignal(sdp: ans, candidate: nil) completion(nil) } } } case "answer": guard let sdpObj = obj["sdp"] as? [String: Any], let sdpStr = sdpObj["sdp"] as? String else { completion(EngineError("answer missing sdp")); return } let desc = RTCSessionDescription(type: .answer, sdp: sdpStr) pc.setRemoteDescription(desc) { err in completion(err) } case "ice": // 空候选 / end-of-candidates:忽略(pion 同样跳过 nil candidate)。 guard let candObj = obj["candidate"] as? [String: Any], let candStr = candObj["candidate"] as? String, !candStr.isEmpty else { completion(nil); return } let mid = candObj["sdpMid"] as? String let mline = (candObj["sdpMLineIndex"] as? NSNumber)?.int32Value ?? 0 let cand = RTCIceCandidate(sdp: candStr, sdpMLineIndex: mline, sdpMid: mid) pc.add(cand) { err in completion(err) } default: completion(nil) } } // MARK: 出站信令 // 与 session.go emitSignal 同形:有 candidate 即 ice,否则取 sdp 自身类型(offer/answer)。 // JSON 形状与 web p2p.ts 的 SignalPayload / pion 的 signalPayload 逐字段对齐(键序无关,两侧皆 parse)。 private func emitSignal(sdp: RTCSessionDescription?, candidate: RTCIceCandidate?) { var payload: [String: Any] = [:] if let candidate { payload["type"] = "ice" // RTCIceCandidate(ObjC)不暴露 usernameFragment,故 candidate JSON 仅含 candidate/sdpMid/ // sdpMLineIndex。ufrag 由 m-line 的 ice-ufrag 决定(offer/answer 已带),pion 与浏览器的 // addIceCandidate 对缺省 usernameFragment 皆容忍,互通不受影响。 var c: [String: Any] = [ "candidate": candidate.sdp, "sdpMLineIndex": Int(candidate.sdpMLineIndex), ] c["sdpMid"] = candidate.sdpMid ?? NSNull() payload["candidate"] = c } else if let sdp { let t = Self.sdpTypeString(sdp.type) payload["type"] = t payload["sdp"] = [ "type": t, "sdp": sdp.sdp ] } else { return } guard let data = try? JSONSerialization.data(withJSONObject: payload), let json = String(data: data, encoding: .utf8) else { return } engine?.events?.engineDidEmitSignal(sessionId: sessionId, toPeer: peerName, payloadJSON: json) } // MARK: 发送端 private func startSenderStream() { cLock.lock() if senderStarted { cLock.unlock(); return } senderStarted = true cLock.unlock() streamQueue.async { [weak self] in self?.streamFile() } } // streamFile 把文件按 64KB 分片经数据通道发出,水位回压与桌面 pion / web 引擎一致:缓冲将越 // HIGH 前等抽干到 LOW。发完 done 帧后先抽干本地缓冲、再等接收端 ack 追平总量才宣告 completed。 private func streamFile() { guard let dc else { fail("no data channel"); return } let fh: FileHandle do { fh = try FileHandle(forReadingFrom: URL(fileURLWithPath: filePath)) } catch { fail("open \(filePath): \(error.localizedDescription)"); return } defer { try? fh.close() } if !sendText([ "type": "meta", "name": (filePath as NSString).lastPathComponent, "size": fileSize ]) { fail("send meta"); return } while true { if isDone() { return } let chunk: Data? do { chunk = try fh.read(upToCount: Wire.chunkSize) } catch { fail("read \(filePath): \(error.localizedDescription)"); return } guard let chunk, !chunk.isEmpty else { break } // EOF if dc.bufferedAmount + UInt64(chunk.count) > Wire.highWatermark { waitBufferLow(Wire.lowWatermark) } if isDone() { return } if !dc.sendData(RTCDataBuffer(data: chunk, isBinary: true)) { fail("send chunk failed (channel closed?)"); return } addSent(Int64(chunk.count)) emitProgress() } _ = sendText([ "type": "done" ]) waitDrained() waitAck(fileSize) emitProgressNow() setState("completed") } @discardableResult private func sendText(_ dict: [String: Any]) -> Bool { guard let dc, let data = try? JSONSerialization.data(withJSONObject: dict) else { return false } return dc.sendData(RTCDataBuffer(data: data, isBinary: false)) } // 等本地缓冲落到 threshold 以下:靠 didChangeBufferedAmount 唤醒 bufSem,兼 50ms 轮询兜底。 private func waitBufferLow(_ threshold: UInt64) { let start = Date() var warned = false while true { guard let dc else { return } if dc.bufferedAmount <= threshold { return } if isDone() { return } _ = bufSem.wait(timeout: .now() + 0.05) if !warned, Date().timeIntervalSince(start) > 5 { warned = true log("send stalled >5s: buffered=\(dc.bufferedAmount) sent=\(getSent()) acked=\(getAcked())") } } } // 等本地缓冲彻底抽干(done 帧也已离开本端 SCTP)。 private func waitDrained() { while true { guard let dc else { return } if dc.bufferedAmount == 0 { return } if isDone() { return } _ = bufSem.wait(timeout: .now() + 0.02) } } // 等接收端 ack 追平 total(确认已全收)。超时兜底;无任何 ack 则不空等(对齐 session.go waitAck)。 private func waitAck(_ total: Int64) { if getAcked() == 0 { return } let deadline = Date().addingTimeInterval(Wire.ackCompleteTimeout) while getAcked() < total { if Date() > deadline { return } if isDone() { return } _ = bufSem.wait(timeout: .now() + 0.05) } } // MARK: 接收端 private func bindReceiverDataChannel(_ channel: RTCDataChannel) { dc = channel channel.delegate = self startAckLoop() } // 每 200ms 把最新已收字节回传发送端(节流,避免每片都回与下行争用通道)。 private func startAckLoop() { let timer = DispatchSource.makeTimerSource(queue: streamQueue) timer.schedule(deadline: .now() + .milliseconds(Wire.ackIntervalMs), repeating: .milliseconds(Wire.ackIntervalMs)) var lastSent: Int64 = -1 timer.setEventHandler { [weak self] in guard let self else { return } let r = self.getReceived() if r != lastSent { self.sendAck(r); lastSent = r } } // 若 cleanup 已先发生(极少见:didOpen 与 cancel 同刻),不再挂表,避免泄漏一个活计时器。 stateLock.lock() if cleanedUp { stateLock.unlock(); timer.cancel(); return } ackTimer = timer stateLock.unlock() timer.resume() } private func sendAck(_ bytes: Int64) { guard let dc, dc.readyState == .open, let data = try? JSONSerialization.data(withJSONObject: [ "type": "ack", "bytes": bytes ]) else { return } _ = dc.sendData(RTCDataBuffer(data: data, isBinary: false)) } // 文本控制帧:发送端处理 ack,接收端处理 meta / done。角色守护避免误处理对端不该发的帧。 private func handleControl(_ data: Data) { guard let obj = (try? JSONSerialization.jsonObject(with: data)) as? [String: Any], let type = obj["type"] as? String else { return } switch type { case "meta": if role == .receiver { openOutput(name: obj["name"] as? String ?? "download") } case "done": if role == .receiver { finalizeReceive() } case "ack": if role == .sender, let b = (obj["bytes"] as? NSNumber)?.int64Value { setAcked(b) bufSem.signal() // 唤醒 waitAck / waitDrained emitProgress() } default: break } } private func openOutput(name: String) { outName = name let dir = engine?.downloadDir() ?? NSTemporaryDirectory() let tmp = (dir as NSString).appendingPathComponent(".\(sessionId).part") tmpPath = tmp FileManager.default.createFile(atPath: tmp, contents: nil) let handle = FileHandle(forWritingAtPath: tmp) outLock.lock(); outHandle = handle; outLock.unlock() if handle == nil { fail("create \(tmp)") } } // 直接写盘——无 base64 桥、无 OPFS、无 WebView 内存。dc 委托串行投递,写入天然有序。 private func writeChunk(_ data: Data) { // 持 outLock 整段写:与 cleanup 的关句柄互斥,杜绝写到已关闭 / 已释放的句柄(见 outHandle 注释)。 outLock.lock() if let h = outHandle { do { try h.write(contentsOf: data) } catch { outLock.unlock(); fail("write chunk: \(error.localizedDescription)"); return } } outLock.unlock() _ = addReceived(Int64(data.count)) emitProgress() } // 收到 done=所有分片已到(ordered),receivedBytes 已等于总量。回最终 ack、关文件、改名落到 // 不冲突路径,报告路径与 completed。命名避开 NSObject 已弃用的 finalize()(否则被判为 override)。 private func finalizeReceive() { sendAck(getReceived()) outLock.lock(); let h = outHandle; outHandle = nil; outLock.unlock() if let h { try? h.close() let dir = engine?.downloadDir() ?? NSTemporaryDirectory() let final = Self.nonCollidingPath(dir: dir, name: outName) do { try FileManager.default.moveItem(atPath: tmpPath, toPath: final) } catch { fail("rename to \(final): \(error.localizedDescription)"); return } engine?.events?.engineDidSave(sessionId: sessionId, path: final) } emitProgressNow() setState("completed") } // MARK: 进度 / 状态 / 收尾 // 当前应上报的已传字节:发送端优先用接收端 ack 追平值(无 ack 时回退本地已交付估计), // 接收端用已写盘字节(对齐 session.go progressBytes)。 private func progressBytes() -> Int64 { if role == .sender { let a = getAcked() if a > 0 { return min(fileSize, a) } let buffered = Int64(dc?.bufferedAmount ?? 0) let b = getSent() - buffered return b < 0 ? 0 : b } return getReceived() } // 节流上报(~10Hz):高吞吐下每片一回调把宿主主线程打满(每片过桥 evaluateJavaScript)。 private func emitProgress() { let now = Int64(DispatchTime.now().uptimeNanoseconds) cLock.lock() if now - lastEmitNs < Wire.progressThrottleNs { cLock.unlock(); return } lastEmitNs = now cLock.unlock() engine?.events?.engineDidProgress(sessionId: sessionId, bytes: progressBytes()) } // 无视节流强发一次:收尾确保最终字节数到达宿主(否则末次进度被节流吞掉、文案停在 99%)。 private func emitProgressNow() { cLock.lock(); lastEmitNs = Int64(DispatchTime.now().uptimeNanoseconds); cLock.unlock() engine?.events?.engineDidProgress(sessionId: sessionId, bytes: progressBytes()) } private func setState(_ st: String) { stateLock.lock() // 终态不可被覆盖:completed/failed 之后即便连接随后 closed/failed 也不回退。 if state == st || state == "completed" || state == "failed" { stateLock.unlock(); return } state = st stateLock.unlock() engine?.events?.engineDidChangeState(sessionId: sessionId, state: st) switch st { case "failed", "closed": cleanup() case "completed": // 不立即拆连接:接收端刚发最终 ack;高吞吐下反向 ack 可能未达发送端,此刻关通道会让发送端 // 误判失败。留连接,ackLoop 续发;60s 兜底防泄漏。真正拆除交宿主 SSE 终态驱动的 cancel。 streamQueue.asyncAfter(deadline: .now() + 60) { [weak self] in self?.cleanup() } default: break } } func fail(_ msg: String) { log("fail: \(msg)") setState("failed") } func cleanup() { stateLock.lock() if cleanedUp { stateLock.unlock(); return } cleanedUp = true doneFlag = true let timer = ackTimer // 与 startAckLoop 的赋值经 stateLock 互斥,避免 ackTimer ARC 竞争 ackTimer = nil stateLock.unlock() bufSem.signal() // 唤醒任何在等缓冲 / ack 的发送循环 timer?.cancel() // outHandle 经 outLock 取出再关:与 writeChunk 互斥(writeChunk 持锁期间 cleanup 在此阻塞, // 待其写完再关),杜绝并发 close/write 与 ARC 竞争。 outLock.lock(); let h = outHandle; outHandle = nil; outLock.unlock() if let h { try? h.close() } dc?.close() pc?.close() engine?.remove(peerName, self) } // MARK: 计数器(cLock 守护) private func addSent(_ n: Int64) { cLock.lock(); bytesSent += n; cLock.unlock() } private func getSent() -> Int64 { cLock.lock(); defer { cLock.unlock() }; return bytesSent } private func addReceived(_ n: Int64) -> Int64 { cLock.lock(); receivedBytes += n; let r = receivedBytes; cLock.unlock(); return r } private func getReceived() -> Int64 { cLock.lock(); defer { cLock.unlock() }; return receivedBytes } private func setAcked(_ v: Int64) { cLock.lock(); if v > ackedBytes { ackedBytes = v }; cLock.unlock() } private func getAcked() -> Int64 { cLock.lock(); defer { cLock.unlock() }; return ackedBytes } private func isDone() -> Bool { stateLock.lock(); defer { stateLock.unlock() }; return doneFlag } // MARK: 选中候选对(连接建立后异步取,best-effort,仅供调试面板) private func logSelectedPair() { // 不用 streamQueue:发送端的 streamFile 会长期占用它(背压等待中),用全局队列才能在传输期间 // 实时取到选中候选对(调试面板用),而非等传输结束。 DispatchQueue.global(qos: .utility).asyncAfter(deadline: .now() + 0.3) { [weak self] in self?.querySelectedPair(attempt: 0) } } private func querySelectedPair(attempt: Int) { guard attempt < 12, !isDone(), let pc else { return } pc.statistics { [weak self] report in guard let self else { return } if let pair = Self.extractSelectedPair(report) { self.log("ice pair: local=\(pair.0) remote=\(pair.1)") self.engine?.events?.engineDidSelectIcePair(sessionId: self.sessionId, local: pair.0, remote: pair.1) } else { DispatchQueue.global(qos: .utility).asyncAfter(deadline: .now() + 0.3) { [weak self] in self?.querySelectedPair(attempt: attempt + 1) } } } } // MARK: 静态助手 private static func parseICEServers(_ raw: String) -> [RTCIceServer] { let trimmed = raw.trimmingCharacters(in: .whitespacesAndNewlines) guard !trimmed.isEmpty, let data = trimmed.data(using: .utf8), let arr = (try? JSONSerialization.jsonObject(with: data)) as? [[String: Any]] else { return [] } var out: [RTCIceServer] = [] for s in arr { var urls: [String] = [] if let one = s["urls"] as? String { urls = [one] } else if let many = s["urls"] as? [String] { urls = many } if urls.isEmpty { continue } let username = s["username"] as? String let credential = s["credential"] as? String if let username, let credential, !username.isEmpty { out.append(RTCIceServer(urlStrings: urls, username: username, credential: credential)) } else { out.append(RTCIceServer(urlStrings: urls)) } } return out } // 返回 dir 下不与现有文件冲突的目标路径:name、name (1)、name (2)……(对齐 wire.go nonCollidingPath)。 private static func nonCollidingPath(dir: String, name: String) -> String { let fm = FileManager.default let base = (dir as NSString).appendingPathComponent(name) if !fm.fileExists(atPath: base) { return base } let ext = (name as NSString).pathExtension let stem = (name as NSString).deletingPathExtension var i = 1 while true { let candidate = ext.isEmpty ? "\(stem) (\(i))" : "\(stem) (\(i)).\(ext)" let p = (dir as NSString).appendingPathComponent(candidate) if !fm.fileExists(atPath: p) { return p } i += 1 } } private static func extractSelectedPair(_ report: RTCStatisticsReport) -> (String, String)? { let stats = report.statistics var localId: String? var remoteId: String? for (_, s) in stats where s.type == "candidate-pair" { let v = s.values let succeeded = (v["state"] as? String) == "succeeded" let nominated = (v["nominated"] as? NSNumber)?.boolValue ?? false if succeeded, nominated { localId = v["localCandidateId"] as? String remoteId = v["remoteCandidateId"] as? String break } } guard let localId, let remoteId else { return nil } func label(_ id: String) -> String { guard let c = stats[id] else { return "?" } let t = (c.values["candidateType"] as? String) ?? "?" let proto = (c.values["protocol"] as? String) ?? "?" return "\(t)/\(proto)" } return (label(localId), label(remoteId)) } private static func sdpTypeString(_ t: RTCSdpType) -> String { switch t { case .offer: return "offer" case .answer: return "answer" case .prAnswer: return "pranswer" case .rollback: return "rollback" @unknown default: return "offer" } } private static func iceConnString(_ s: RTCIceConnectionState) -> String { switch s { case .new: return "new" case .checking: return "checking" case .connected: return "connected" case .completed: return "completed" case .failed: return "failed" case .disconnected: return "disconnected" case .closed: return "closed" case .count: return "count" @unknown default: return "?" } } // 从候选 SDP 串提取 "type/proto"(仅日志,best-effort): // candidate: typ ... private static func candTypeString(_ sdp: String) -> String { let parts = sdp.split(separator: " ").map(String.init) var typ = "?" var proto = "?" if parts.count > 2 { proto = parts[2].lowercased() } if let i = parts.firstIndex(of: "typ"), i + 1 < parts.count { typ = parts[i + 1] } return "\(typ)/\(proto)" } private func log(_ msg: String) { engine?.events?.engineDidLog("session \(sessionId) \(msg)") } } // MARK: - RTCPeerConnectionDelegate(回调在 libwebrtc 信令线程,非主线程) extension P2PTransferSession: RTCPeerConnectionDelegate { func peerConnection(_ peerConnection: RTCPeerConnection, didChange stateChanged: RTCSignalingState) {} func peerConnection(_ peerConnection: RTCPeerConnection, didAdd stream: RTCMediaStream) {} func peerConnection(_ peerConnection: RTCPeerConnection, didRemove stream: RTCMediaStream) {} func peerConnectionShouldNegotiate(_ peerConnection: RTCPeerConnection) {} func peerConnection(_ peerConnection: RTCPeerConnection, didChange newState: RTCIceConnectionState) { log("ice: \(Self.iceConnString(newState))") } func peerConnection(_ peerConnection: RTCPeerConnection, didChange newState: RTCIceGatheringState) { if newState == .complete { log("local gathering done") } } func peerConnection(_ peerConnection: RTCPeerConnection, didGenerate candidate: RTCIceCandidate) { log("cand local: \(Self.candTypeString(candidate.sdp))") emitSignal(sdp: nil, candidate: candidate) } func peerConnection(_ peerConnection: RTCPeerConnection, didRemove candidates: [RTCIceCandidate]) {} func peerConnection(_ peerConnection: RTCPeerConnection, didOpen dataChannel: RTCDataChannel) { if role == .receiver { bindReceiverDataChannel(dataChannel) } } // 组合 ice+dtls 连接态(等价 session.go OnConnectionStateChange):驱动 connected/failed/closed。 func peerConnection(_ peerConnection: RTCPeerConnection, didChange newState: RTCPeerConnectionState) { switch newState { case .connected: setState("connected") logSelectedPair() case .failed: fail("peerconnection failed") case .closed: setState("closed") default: break } } } // MARK: - RTCDataChannelDelegate(回调在 libwebrtc 信令线程) extension P2PTransferSession: RTCDataChannelDelegate { func dataChannelDidChangeState(_ dataChannel: RTCDataChannel) { // 发送端:通道就绪即开始 streamFile(等价 session.go dc.OnOpen → go streamFile)。 if role == .sender, dataChannel.readyState == .open { startSenderStream() } } func dataChannel(_ dataChannel: RTCDataChannel, didReceiveMessageWith buffer: RTCDataBuffer) { if buffer.isBinary { writeChunk(buffer.data) } else { handleControl(buffer.data) } } func dataChannel(_ dataChannel: RTCDataChannel, didChangeBufferedAmount amount: UInt64) { if amount <= Wire.lowWatermark { bufSem.signal() } } }