diff --git a/desktop/.gitignore b/desktop/.gitignore index 129d522..2c09a96 100644 --- a/desktop/.gitignore +++ b/desktop/.gitignore @@ -1,3 +1,4 @@ build/bin node_modules frontend/dist +/cdrop-desktop diff --git a/desktop/NATIVE-TRANSFER.md b/desktop/NATIVE-TRANSFER.md new file mode 100644 index 0000000..f0c8fc5 --- /dev/null +++ b/desktop/NATIVE-TRANSFER.md @@ -0,0 +1,179 @@ +# cdrop 桌面原生数据面(Option A)设计 + +> 目标:把桌面 **P2P 数据面**(WebRTC DataChannel 收发)从 WebView 的 JS 移到 Go 进程(`pion/webrtc`),**留在 Wails、不上 Tauri**。 +> 关联:`desktop/PLAN.md`(桌面总计划)、`desktop/RESEARCH.md`(调研)、`web/src/features/transfer/{transfer,p2p,relay,incomingSink,source}.ts`、吞吐专项诊断(task #38)。 +> 状态:设计稿(2026-06-27)。Phase 0(止血)与 Phase 1(落地)见末尾阶段表。 + +--- + +## 0. 为什么是 Go/pion,而不是 Tauri + +吞吐专项(#38)定位出**两个根因**,都活在 WebView 的 WebRTC 层: + +- **B2 渲染器节流**:Chromium/WebView2 的 SCTP 处理线程活在**渲染器进程**内,窗口隐藏/遮挡/最小化触发 renderer backgrounding + 定时器节流 → SCTP 被饿。Windows host↔host 直连仍 100KB/s 即此。 +- **B1 接收窗口**:WebKit 的 usrsctp 接收窗口写死 256KB(JS 改不动),`吞吐 ≈ rwnd / RTT`,高 RTT 路径被掐。 + +**Tauri 用的是同款系统 WebView(WebView2 / WKWebView),上述两条一个不变**——换 Tauri 是把整个 `desktop/platform/` Go 树重写成 Rust、零吞吐收益。真正的杠杆是「把数据面移出 WebView 跑原生」,而这在**现有 Go 栈用 `pion/webrtc` 即可**,无须 Tauri。 + +**原生 Go/pion 一次性收口:** + +| 根因 | 留 WebView | 原生 Go/pion | 机理 | +|---|---|---|---| +| B2 渲染节流 | 旗标缓解、仍受启发式摆布 | **彻底消失** | WebRTC 不在渲染器里,窗口状态与网络栈无关——「网络不被渲染阻塞」的终极形态 | +| B1 桌面作接收端 | WebKit 写死 256KB | **可配** | `SettingEngine.SetSCTPMaxReceiveBufferSize` 设 1-16MB | +| mac srflx 非 host | WebKit mDNS 隐私替换 | **消失** | pion 默认用真实 IP 的 host 候选,不做 `.local` 混淆 | +| base64 过桥 | 收发都过 base64 桥 | **消失** | Go 直接读/写磁盘文件,无 OPFS / 混合 sink | +| 桌面↔桌面 | 双重受限 | **全速** | 两端原生大 rwnd + 直连 | + +--- + +## 1. 范围 + +**移到 Go:** 仅 **P2P 数据面**(pion peerconnection / datachannel 收发、文件磁盘 I/O)**及其信令**(offer/answer/ICE)。 + +**留在 JS(单实现、不动):** 全部**编排**——`transfer.ts` 的 `/initiate`、30s 中继回退看门狗、离线对端 presence 等待、relay 收发、状态机 POST、取消、store/UI。relay 非吞吐瓶颈(~3MB/s,且本就过 Go 代理),Phase 1 不碰。 + +**残差(Phase 1 不解):** 桌面 → WebKit 接收端(iOS/Safari)在**高 RTT** 路径仍 `256KB / RTT`——留 Phase 2 striping。同内网低 RTT 路径下 256KB 非约束,不受影响。 + +--- + +## 2. 接缝:`p2pBackend` 抽象(核心,使 `transfer.ts` 零改动) + +现状调用链:`transfer.ts` / `hub.ts` → `p2pStartOutgoing(sessionId, peer, src) : P2PSession`、`p2pStartIncoming(sessionId, peer) : P2PSession`、`p2pHandleSignal(from, payload)`、`p2pCleanup(peer)`。 + +引入同形后端接口: + +``` +interface P2pBackend + startOutgoing(sessionId, peerName, src) → P2PSession // P2PSession: onState / onProgress / cancel,与现有一致 + startIncoming(sessionId, peerName) → P2PSession + handleSignal(fromPeer, payload) → void + cleanup(peerName) → void + +selectBackend() → isDesktop() ? goBridgeBackend : jsWebrtcBackend // web / iOS 走 jsWebrtcBackend(现有 p2p.ts) +``` + +`transfer.ts` 的编排(看门狗 / presence / relay / state POST / cancel)**全部不变**——它只面对 `P2PSession` 事件接口;后端是 JS-WebRTC 还是 Go-bridge 对它透明。 + +--- + +## 3. Go 侧:`p2pengine` 包(pion) + +``` +Sender(sessionId, peerName, filePath, iceServers): + pc ← pion.NewPeerConnection(SettingEngine{ rwnd 大; host 候选真 IP; iceServers }) + dc ← pc.CreateDataChannel("cdrop-file", ordered=true) + onOpen: + send JSON {type:meta, name, size, sha256?} // 与 p2p.ts 同帧 + for chunk in readFileByPath(filePath, CHUNK): // Go 直接磁盘读,无桥、无 base64 + backpressure: 等 dc.BufferedAmount 落水位 + dc.Send(chunk) // 二进制帧 + send JSON {type:done} + 等 receiver ack 追平 size → 报 completed + onMessage(ack): 更新 ackedBytes → emit progress + POST /api/transfer/{id}/p2p(活跃); 完成不 POST /done(接收端权威) + +Receiver(sessionId, peerName, iceServers): + pc ← pion.NewPeerConnection(SettingEngine{ rwnd 大; ... }) + onDataChannel(dc): + onMessage(meta): 打开下载目录文件句柄(复用 download.go) + onMessage(chunk): 直接写盘; 节流 200ms 回 ack {type:ack, bytes} + onMessage(done): 立即回最终 ack; 关闭文件; POST /api/transfer/{id}/done + +信令(Go own):pc.OnICECandidate → POST /api/hub/signal(带 a.token) + inbound offer/answer/ice ← JS hub 转发进来(见 §4) +ICE creds:Go 拉 /api/calls/credentials(带 a.token),与 JS 同源 +``` + +**关键配置:** `SetSCTPMaxReceiveBufferSize`(大 rwnd);host 候选真 IP(pion 默认,无 mDNS);复用现有 Cloudflare TURN。 + +--- + +## 4. 桥接面(Wails Go↔JS)与「单 SSE」原则 + +**JS → Go(bind 方法):** +- `StartOutgoing(sessionId, peerName, filePath)`、`StartIncoming(sessionId, peerName)` +- `HandleSignal(sessionId, fromPeer, payloadJSON)` ← JS hub 把入站信令转发给 Go +- `Cancel(sessionId)` ← 取消 / 中继回退时 JS 通知 Go 拆 pion + +**Go → JS(Wails events):** `transfer:progress {sessionId, bytes}`、`transfer:state {sessionId, state}`。 + +**单 SSE 原则(重要):** 设备只有**一条** SSE(JS `hub.ts`,承载 presence + 信令 + 唤醒)。**不**让 Go 另开 SSE——第二条会污染 presence/在线态。故入站信令统一由 JS hub 收下,对 Go-owned 会话转发给 Go;出站信令由 Go 直接 POST(Go 有 `a.token` + `oauth.go` 刷新,本就是 API 调用的正确归属)。 + +--- + +## 5. 发送取文件:改用 Wails 原生路径 + +浏览器 `` / drag-drop 给的是沙箱 `File`(**无磁盘路径**、受 256MB 内存约束)。Go 要直接读盘,需**绝对路径**: + +- 桌面发送改用 **Wails 原生文件对话框 / `OnFileDrop`**(返回绝对路径)→ 传路径给 Go。 +- `name`/`size` 由 Go `stat` 出,回给 JS 供 `/initiate` 与 store 记录。 +- 这条是桌面前端的局部改动(`isDesktop()` 分支),收益:彻底脱离 WebView 内存、无桥。 + +--- + +## 6. 协议对齐清单(JS 引擎 ↔ Go 引擎必须逐项一致) + +> 双实现的唯一契约是**线协议**。任何一项漂移都会让桌面↔web/iOS 互通断裂。建议抽出共享测试向量(meta/chunk/done/ack 的字节级样例)双侧回归。 + +- **信令**:`POST /api/hub/signal {to, payload:{type:"offer"|"answer"|"ice", sdp?, candidate?}}`;入站经 SSE。 +- **DataChannel**:名 `"cdrop-file"`,`ordered:true`。 +- **帧**:`meta` JSON `{type,name,size,sha256?}` / 二进制 chunk / `done` JSON / `ack` JSON `{type,bytes}`(接收端节流 200ms 回传累计已收字节,单调取大)。 +- **状态机**:`/api/transfer/{initiate,p2p,done,fail,cancel,fallback}`;**完成由接收端 POST `/done`**(权威),发送端不重复。 +- **中继回退语义**:JS 看门狗 30s 未 `connected` → `POST /fallback` 并对 Go `Cancel(session)`;Go 拆 pion,relay 由 JS 接管(数据面回到 JS relay,符合「relay 留 JS」)。 + +--- + +## 7. iOS 数据面:经 gomobile 与桌面共用同一份 Go 引擎(U1,POC 闸) + +**决策(2026-06-27,研究后翻案)**:iOS 不再「维持现状 + 仅靠 striping」,而是**把数据面也迁到 Go/pion,经 `gomobile bind` 与桌面共用同一份引擎**——原生侧只一份实现(web 仍 JS,共 2 份引擎),iOS 彻底摆脱 WebKit 256KB。依据用户方针「多端统一 + 各情况性能」,且桌面已定 pion——iOS 同走 pion 即与桌面**同引擎、同特性、同协议**,统一性最大。**该决策以一个 2 天 POC 为闸**(见 §8 Phase 2)。 + +**研究结论(U1 = gomobile + pion):** +- **可行,但不直接 bind pion**(pion 公开 API 含 gobind 不支持的类型,pion#1111)。正解=写一层 gobind-clean 的 `engine` 薄包装内部引 pion,再 `gomobile bind ./engine` 出 `.xcframework`。pion 维护者本人推荐 gomobile(discussion#1746)。 +- **FFI 边界对 cdrop 恰好够用**:API(`StartOutgoing/HandleSignal/Cancel` 全字符串 + 回调接口 `OnProgress(int64)/OnState/OnSignal(json)`)完全合法;**文件 chunk 的 `[]byte` 不过边界**——Go 自读盘经 pion 内部发送,只 int + JSON 串过桥。 +- **rwnd 可配实锤**:`SetSCTPMaxReceiveBufferSize(4MB)` 广告 4MB 接收窗,碾过 256KB;另有 `SetSCTPMinCwnd / SetSCTPCwndCAStep / EnableSCTPZeroChecksum` 等拥塞/延迟旋钮。 +- **关 mDNS + 真 host 候选**(`SetAnsweringMDNSEnabled(false)`):从根上避开 srflx 退化;首次 LAN 连接弹一次本地网络权限(需 `NSLocalNetworkUsageDescription`)。 +- **体积 ~9-15MB(压缩),小于 libwebrtc**;后台挂起需 `beginBackgroundTask`(iOS 通病、非缺陷)。 +- **成本 ~1 周**;成熟度 LOW-MEDIUM(无已知公开 pion+gomobile 生产案例——主要风险)。 + +**架构对称性(统一的真正收益)**:iOS 变成与桌面同构—— +- 桌面:Wails 壳(JS 编排 + UI)+ Go/pion 数据面(原生进程)。 +- iOS:SwiftUI 壳(WKWebView 跑 JS 编排/hub/presence/messages + UI)+ Go/pion 数据面(gomobile xcframework)。 +- 二者共用同一 `p2pBackend` 抽象、同一桥接语义(单 SSE 留 JS hub、信令转发给 Go)、同一 wire protocol、同一 striping 实现。差别仅在原生壳与桥的 binding(Wails bind vs gomobile + WKScriptMessageHandler)。 + +**U1 vs U2(备选)对照:** + +| | U1:gomobile + pion(iOS 与桌面同引擎) | U2:iOS 用 libwebrtc.framework | +|---|---|---| +| 统一 | **2 份引擎**(JS + Go),原生侧一份实现 | 3 份(JS-browser / Go-pion / iOS-libwebrtc) | +| iOS rwnd | pion 4MB+(可配) | DcSCTP **5MB 默认**(不可经 ObjC 配,但默认已够) | +| 高 RTT 吞吐 | 较弱(pion/sctp#218,~17-45Mbps@100-300ms,仍 10-20× 于 256KB);可经 cwnd 旋钮 + striping 补 | **更强**(libwebrtc 成熟拥塞控制) | +| LAN 吞吐 | 优(177-218Mbps 实测) | 优(原生,去浏览器 IPC 开销) | +| 体积 | ~9-15MB | ~12-18MB(设备切片更大) | +| 成熟度/风险 | LOW-MEDIUM(gomobile 无生产先例、Xcode 版本偶断) | 成熟(stasel/livekit 预编包),但**多一份栈**、社区维护需 pin 版本 | +| 互通 | pion↔browser 经 pion CI 持续验证;iOS-gomobile 网络栈待 POC 验 | iOS-libwebrtc↔pion **无文档**,需冒烟测(pion#2288 报抖动) | + +**裁决:选 U1,以 POC 为闸。** 理由:桌面已 pion → iOS 同走 pion 保持引擎/特性/协议一致、避免第三份栈,最契合「统一」;性能在 cdrop 的局域网主场景两者皆优;高 RTT 的差距是 **pion 全局问题**(桌面同样吃),单点解决(cwnd 旋钮 / striping / 上游修 #218 / 就近 TURN)即全端受益,且 pion 高 RTT 仍 10-20× 于今日 WebKit。**以可补的高 RTT 残差,换最大化的统一收益,符合用户方针。** + +**统一在协议 + 实现(双层)**:web 无原生选项故 JS 引擎不可消;但桌面与 iOS 经 U1 收敛到**同一份 Go 引擎**。三端共享**同一 wire protocol**(§6),striping 作为协议扩展在「JS 引擎(web)」与「Go 引擎(桌面+iOS)」两处实现。 + +--- + +## 8. 阶段 + +| 阶段 | 内容 | 验收 | +|---|---|---| +| **Phase 0(止血,临时)** | C:`WEBVIEW2_ADDITIONAL_BROWSER_ARGUMENTS` 注入反节流旗标(`--disable-renderer-backgrounding` 等) | Windows 隐藏/遮挡窗口传输不掉速 | +| **Phase 1(落地 A·桌面)** | Go `p2pengine`(pion,单关联)+ `p2pBackend` 抽象 + 桥接 + 原生取文件路径;relay/编排留 JS | 桌面↔桌面/Chrome/iOS-LAN **双向快**;隐藏窗口传输满速;去 base64 桥;mac 落 host↔host | +| **Phase 2(iOS POC 闸,~2 天)** | 最小 `engine` 包装 → `gomobile bind -target ios`(用最新 `golang.org/x/mobile`)→ Swift 调 `StartOutgoing` | ① 编译出 xcframework;② iOS 上 Go 开 UDP 拿到 **host 候选** + 本地网络权限弹窗;③ 真跑 `iOS(pion)↔桌面(pion)` 与 `iOS(pion)↔浏览器` 传输达预期。**过→锁 U1;撞硬阻塞→退 U2** | +| **Phase 3(落地 U1·iOS)** | iOS 数据面接入共享 Go 引擎(gomobile xcframework + iOS 侧 `p2pBackend` binding);JS 编排/hub/presence 留 WKWebView | iOS 摆脱 256KB;iOS↔* 高 RTT 大幅改善;与桌面**同引擎** | +| **Phase 4(按需)** | 共享协议 **N 关联 striping**(JS 引擎 for web;Go 引擎 for 桌面+iOS) | 解跨 NAT 高 RTT 残差(全端) | + +--- + +## 9. 风险与缓解 + +- **pion ↔ 浏览器互通**:标准 SDP/DTLS/SCTP,pion 主用例,低风险;**需验** pion DataChannel 在大 rwnd 下的实测吞吐。 +- **双实现协议漂移**:靠 §6 对齐清单 + 共享字节级测试向量双侧回归。 +- **Go 调 API 鉴权**:Go 已持 `a.token` + `oauth.go` 刷新,本就是 API 调用正确归属;信令/状态机 POST 直接带 Authorization。 +- **前端发送路径分叉**:`isDesktop()` 改用 Wails 原生路径——局部、有先例(`net/desktop` 已有 `isDesktop()` 分流)。 diff --git a/desktop/app.go b/desktop/app.go index 3ef4c3a..7103879 100644 --- a/desktop/app.go +++ b/desktop/app.go @@ -2,14 +2,18 @@ package main import ( "context" + "encoding/base64" "errors" "fmt" + "io" "os" + "path/filepath" "github.com/wailsapp/wails/v2/pkg/menu" "github.com/wailsapp/wails/v2/pkg/menu/keys" "github.com/wailsapp/wails/v2/pkg/runtime" + "cdrop-desktop/engine" "cdrop-desktop/platform" ) @@ -24,6 +28,11 @@ type App struct { clipSource platform.Source clipMonitor *platform.Monitor + + // transfer 是原生 P2P 数据面(pion):取代 WebView 内 JS WebRTC,逃离 WebKit/WebView2 + // 的 256KB rwnd 与渲染器节流。HTTP(信令 / 状态机)仍由 WebView 侧 JS 经事件桥发起, + // 引擎只做纯 WebRTC + 文件 I/O。见 desktop/NATIVE-TRANSFER.md。 + transfer *engine.Engine } // NewApp creates a new App application struct. @@ -53,6 +62,8 @@ func (a *App) startup(ctx context.Context) { if platform.IsLaunchAtLoginEnabled() { _ = platform.SetLaunchAtLogin(true) } + // 原生传输引擎:落地目录取当前配置(设置页改目录时经 SaveSettings 同步到引擎)。 + a.transfer = engine.New(engine.Config{DownloadDir: platform.ResolveDownloadDir()}, &transferEvents{app: a}) a.startClipboardSync(ctx) // 触发 macOS 本地网络权限(macOS 15+ 隐私门):使本进程 WKWebView 的 WebRTC 能收集 host / // mDNS 候选、同内网走直连而非 prflx↔prflx 慢路径(见 platform.TriggerLocalNetwork)。其他平台空实现。 @@ -325,6 +336,10 @@ func (a *App) SaveSettings(cfg platform.DesktopConfig) error { if err := platform.SaveConfig(cur); err != nil { return err } + // 让引擎的落地目录跟随设置变化(接收端直接写盘)。 + if a.transfer != nil { + a.transfer.SetDownloadDir(platform.ResolveDownloadDir()) + } return platform.SetLaunchAtLogin(cur.LaunchAtLogin) } @@ -422,3 +437,103 @@ func (a *App) LoggedIn() bool { func (a *App) Greet(name string) string { return fmt.Sprintf("Hello %s, It's show time!", name) } + +// --- 原生 P2P 数据面桥(绑定到 WebView) --- +// +// WebView 侧的 p2pBackend(Go-bridge 后端)调用这四个方法发起 / 接收 / 喂信令 / 取消; +// 引擎经 transferEvents 反向 EventsEmit 进度 / 状态 / 出站信令 / 落盘路径。HTTP 全留 JS: +// 出站信令由 JS POST /api/hub/signal,入站信令由 JS hub 转发进 P2PHandleSignal,状态机 +// POST(/p2p、/done、/fail)由 JS 据 p2p:state 事件发起。iceServersJSON 由 JS 透传 +// (与 web 同源,来自 /api/calls/credentials)。 + +// P2PStartOutgoing 作为发送端发起:filePath 为本机绝对路径(经 PickFileForSend 获得)。 +func (a *App) P2PStartOutgoing(sessionID, peerName, filePath, iceServersJSON string) error { + return a.transfer.StartOutgoing(sessionID, peerName, filePath, iceServersJSON) +} + +// P2PStartIncoming 预备接收端,等对端 offer。 +func (a *App) P2PStartIncoming(sessionID, peerName, iceServersJSON string) error { + return a.transfer.StartIncoming(sessionID, peerName, iceServersJSON) +} + +// P2PHandleSignal 把 JS hub 从 SSE 收到的入站信令喂给引擎(按对端设备名路由)。 +func (a *App) P2PHandleSignal(fromPeer, payloadJSON string) error { + return a.transfer.HandleSignal(fromPeer, payloadJSON) +} + +// P2PCancel 拆除会话(取消 / 中继回退)。 +func (a *App) P2PCancel(sessionID string) { + a.transfer.Cancel(sessionID) +} + +// PickedFile 是原生文件选择的结果,供发送端走原生数据面(Go 直接按路径读盘,不经浏览器 +// File 的沙箱无路径 + 内存约束)。 +type PickedFile struct { + Path string `json:"path"` + Name string `json:"name"` + Size int64 `json:"size"` +} + +// ReadFileSlice 读 path 的 [start,end) 字节并 base64 返回。仅供中继回退时 WebView 侧的 +// FileSource.slice 用(原生 P2P 由引擎直接按路径读盘、不走此桥);故 base64 开销只在罕见 +// 的 relay 路径上,不影响 P2P 主路径。镜像 iOS 的 readOutgoingSlice。 +func (a *App) ReadFileSlice(path string, start, end int64) (string, error) { + if start < 0 || end < start { + return "", fmt.Errorf("bad range [%d,%d)", start, end) + } + f, err := os.Open(path) + if err != nil { + return "", err + } + defer f.Close() + buf := make([]byte, end-start) + n, err := f.ReadAt(buf, start) + if err != nil && err != io.EOF { + return "", err + } + return base64.StdEncoding.EncodeToString(buf[:n]), nil +} + +// PickFileForSend 打开原生文件对话框,返回选中文件的绝对路径 + 名 + 大小;用户取消返回 nil。 +func (a *App) PickFileForSend(title string) (*PickedFile, error) { + path, err := runtime.OpenFileDialog(a.ctx, runtime.OpenDialogOptions{Title: title}) + if err != nil { + return nil, err + } + if path == "" { + return nil, nil // 用户取消 + } + fi, err := os.Stat(path) + if err != nil { + return nil, err + } + return &PickedFile{Path: path, Name: filepath.Base(path), Size: fi.Size()}, nil +} + +// transferEvents 把引擎回调转成 Wails 事件,供 WebView 侧 p2pBackend 监听。回调可能在 pion +// 任意 goroutine 触发;Wails 的 EventsEmit / Log* 可跨线程安全调用。 +type transferEvents struct{ app *App } + +func (e *transferEvents) OnProgress(sessionID string, bytes int64) { + runtime.EventsEmit(e.app.ctx, "p2p:progress", map[string]any{"sessionId": sessionID, "bytes": bytes}) +} + +func (e *transferEvents) OnState(sessionID, state string) { + runtime.EventsEmit(e.app.ctx, "p2p:state", map[string]any{"sessionId": sessionID, "state": state}) +} + +func (e *transferEvents) OnSignal(sessionID, toPeer, payloadJSON string) { + runtime.EventsEmit(e.app.ctx, "p2p:signal", map[string]any{"sessionId": sessionID, "to": toPeer, "payload": payloadJSON}) +} + +func (e *transferEvents) OnSaved(sessionID, path string) { + runtime.EventsEmit(e.app.ctx, "p2p:saved", map[string]any{"sessionId": sessionID, "path": path}) +} + +func (e *transferEvents) OnIcePair(sessionID, local, remote string) { + runtime.EventsEmit(e.app.ctx, "p2p:icepair", map[string]any{"sessionId": sessionID, "local": local, "remote": remote}) +} + +func (e *transferEvents) OnLog(line string) { + runtime.LogInfof(e.app.ctx, "p2p: %s", line) +} diff --git a/desktop/engine/engine.go b/desktop/engine/engine.go new file mode 100644 index 0000000..e92b2b0 --- /dev/null +++ b/desktop/engine/engine.go @@ -0,0 +1,169 @@ +package engine + +import ( + "fmt" + "os" + "sync" + + "github.com/pion/webrtc/v4" +) + +// Events 由宿主实现,接收引擎的回调。所有方法签名仅用 string / int64——刻意保持 +// gobind-legal,使同一接口将来经 gomobile 暴露给 iOS Swift 时无需改形。 +// +// 回调可能在 pion 的任意 goroutine 上触发;宿主侧的实现需自行保证线程安全 +// (桌面 Wails 的 EventsEmit、iOS 的 DispatchQueue.main 均可跨线程调用)。 +type Events interface { + // OnProgress 报告已传字节:发送端=接收端 ack 追平的字节(无 ack 时回退本地已交付估计), + // 接收端=已写盘字节。 + OnProgress(sessionID string, bytes int64) + // OnState 报告会话状态:connected / completed / failed / closed。宿主据此驱动 store + // 与状态机 POST(/p2p 于 connected、/done 于接收端 completed、/fail 于 failed)。 + OnState(sessionID string, state string) + // OnSignal 交出一条出站信令(offer/answer/ice),宿主负责 POST /api/hub/signal{to,payload}。 + OnSignal(sessionID string, toPeer string, payloadJSON string) + // OnSaved 报告接收端落盘完成的绝对路径,供宿主提示用户。 + OnSaved(sessionID string, path string) + // OnIcePair 报告连接建立后选中的候选对(host/srflx/prflx/relay,含协议,形如 "host/udp"), + // 供宿主在调试面板展示「实际通过哪条路径连上」。 + OnIcePair(sessionID string, local string, remote string) + // OnLog 输出诊断信息。 + OnLog(line string) +} + +// Config 由宿主在构造时提供。DownloadDir 是接收文件的落地目录(桌面 ResolveDownloadDir / +// iOS 应用沙盒)——引擎直接写盘,不再经 base64 桥。 +type Config struct { + DownloadDir string +} + +// Engine 持有所有活跃会话。会话按对端设备名(peer)索引,与 web 引擎 p2pSessions 的键一致: +// 信令仅携带 from/to 设备名、不带 sessionID,故入站信令只能按对端路由(同一对端同一时刻 +// 至多一条会话——沿用 web 端的同等约束)。 +type Engine struct { + cfg Config + ev Events + mu sync.Mutex + byPeer map[string]*session +} + +// New 构造引擎。ev 不可为 nil。 +func New(cfg Config, ev Events) *Engine { + return &Engine{cfg: cfg, ev: ev, byPeer: make(map[string]*session)} +} + +// SetDownloadDir 更新落地目录(设置页改下载目录时调用)。 +func (e *Engine) SetDownloadDir(dir string) { + e.mu.Lock() + e.cfg.DownloadDir = dir + e.mu.Unlock() +} + +func (e *Engine) downloadDir() string { + e.mu.Lock() + defer e.mu.Unlock() + return e.cfg.DownloadDir +} + +// StartOutgoing 作为发送端发起:建 PeerConnection + DataChannel,发 offer,通道就绪后 +// 把 filePath 的内容按线协议流式发出。filePath 为本机绝对路径(桌面经原生文件选择获得)。 +func (e *Engine) StartOutgoing(sessionID, peerName, filePath, iceServersJSON string) error { + fi, err := os.Stat(filePath) + if err != nil { + return fmt.Errorf("stat %s: %w", filePath, err) + } + if fi.IsDir() { + return fmt.Errorf("%s is a directory", filePath) + } + pc, err := newPeerConnection(iceServersJSON) + if err != nil { + return err + } + s := newSession(e, sessionID, peerName, roleSender, pc) + s.filePath = filePath + s.fileSize = fi.Size() + e.put(peerName, s) + s.wireCommon() + + ordered := true + dc, err := pc.CreateDataChannel(channelName, &webrtc.DataChannelInit{Ordered: &ordered}) + if err != nil { + s.fail(fmt.Errorf("create datachannel: %w", err)) + return err + } + s.bindSenderDC(dc) + + offer, err := pc.CreateOffer(nil) + if err != nil { + s.fail(fmt.Errorf("create offer: %w", err)) + return err + } + if err := pc.SetLocalDescription(offer); err != nil { + s.fail(fmt.Errorf("set local description: %w", err)) + return err + } + s.emitSignal(peerName, &offer, nil) + return nil +} + +// StartIncoming 作为接收端预备:建 PeerConnection 等对端 offer。收到 DataChannel 后按 +// 线协议接收并直接落盘。须在对端 offer 到达前调用(宿主在 transfer:incoming 时预备)。 +func (e *Engine) StartIncoming(sessionID, peerName, iceServersJSON string) error { + pc, err := newPeerConnection(iceServersJSON) + if err != nil { + return err + } + s := newSession(e, sessionID, peerName, roleReceiver, pc) + e.put(peerName, s) + s.wireCommon() + pc.OnDataChannel(func(dc *webrtc.DataChannel) { + s.bindReceiverDC(dc) + }) + return nil +} + +// HandleSignal 把宿主从 SSE 收到的入站信令喂给对应会话(按对端设备名路由)。 +func (e *Engine) HandleSignal(fromPeer, payloadJSON string) error { + e.mu.Lock() + s := e.byPeer[fromPeer] + e.mu.Unlock() + if s == nil { + return fmt.Errorf("no active session for peer %q", fromPeer) + } + return s.handleSignal([]byte(payloadJSON)) +} + +// Cancel 拆除指定会话(取消 / 中继回退时由宿主调用)。静默拆除——取消是宿主驱动的, +// 不再回 OnState 以免与宿主自身的终态处理重入。 +func (e *Engine) Cancel(sessionID string) { + e.mu.Lock() + var target *session + for _, s := range e.byPeer { + if s.id == sessionID { + target = s + break + } + } + e.mu.Unlock() + if target != nil { + target.cleanup() + } +} + +func (e *Engine) put(peer string, s *session) { + e.mu.Lock() + // 同一对端若已有会话,先拆旧的(沿用 web 端按对端覆盖的语义)。 + if old := e.byPeer[peer]; old != nil && old != s { + go old.cleanup() + } + e.byPeer[peer] = s + e.mu.Unlock() +} + +func (e *Engine) remove(peer string, s *session) { + e.mu.Lock() + if e.byPeer[peer] == s { + delete(e.byPeer, peer) + } + e.mu.Unlock() +} diff --git a/desktop/engine/engine_test.go b/desktop/engine/engine_test.go new file mode 100644 index 0000000..e4c2ae0 --- /dev/null +++ b/desktop/engine/engine_test.go @@ -0,0 +1,146 @@ +package engine + +import ( + "bytes" + "os" + "path/filepath" + "testing" + "time" +) + +// sigQueue 是测试里每个引擎的入站信令 FIFO:单消费者保序处理,确保 offer/answer 先于 +// 其后 trickle 的 ICE 候选被处理(pion 的 AddICECandidate 在 remote description 未设时报错)。 +type sigQueue struct { + ch chan [2]string // [fromPeer, payloadJSON] + eng *Engine +} + +func newSigQueue() *sigQueue { return &sigQueue{ch: make(chan [2]string, 256)} } + +func (q *sigQueue) run(t *testing.T, done <-chan struct{}) { + go func() { + for { + select { + case <-done: + return + case m := <-q.ch: + if err := q.eng.HandleSignal(m[0], m[1]); err != nil { + t.Logf("HandleSignal(%s) err: %v", m[0], err) + } + } + } + }() +} + +// peerEvents 把一端的出站信令投到对端队列,并暴露 saved/failed 通道供断言。 +type peerEvents struct { + t *testing.T + selfName string + out *sigQueue // 对端队列 + saved chan string + failed chan string +} + +func (p *peerEvents) OnProgress(string, int64) {} +func (p *peerEvents) OnState(_ string, st string) { + if st == "failed" { + select { + case p.failed <- st: + default: + } + } +} +func (p *peerEvents) OnSignal(_ string, _ string, payloadJSON string) { + // 对端看到的 fromPeer 即本端名字。 + p.out.ch <- [2]string{p.selfName, payloadJSON} +} +func (p *peerEvents) OnSaved(_ string, path string) { + select { + case p.saved <- path: + default: + } +} +func (p *peerEvents) OnIcePair(_ string, _ string, _ string) {} +func (p *peerEvents) OnLog(line string) { p.t.Log(line) } + +// TestEngineTransfer 端到端:两个进程内 pion peer(A 发 / B 收)经 host 候选直连,传一个 +// 3MB 已知内容的文件,断言落盘字节与源一致。覆盖线协议(meta/chunk/done/ack)、水位回压、 +// 接收端落盘改名。 +func TestEngineTransfer(t *testing.T) { + dir := t.TempDir() + srcPath := filepath.Join(dir, "payload.bin") + dlDir := filepath.Join(dir, "downloads") + if err := os.MkdirAll(dlDir, 0o755); err != nil { + t.Fatal(err) + } + + // 3MB 确定性内容:触发多片 + 至少一次回压窗口。 + const size = 3 * 1024 * 1024 + want := make([]byte, size) + for i := range want { + want[i] = byte(i*7 + 3) + } + if err := os.WriteFile(srcPath, want, 0o644); err != nil { + t.Fatal(err) + } + + qA, qB := newSigQueue(), newSigQueue() + evA := &peerEvents{t: t, selfName: "A", out: qB, saved: make(chan string, 1), failed: make(chan string, 1)} + evB := &peerEvents{t: t, selfName: "B", out: qA, saved: make(chan string, 1), failed: make(chan string, 1)} + + engA := New(Config{DownloadDir: dlDir}, evA) // 发送端(不落盘) + engB := New(Config{DownloadDir: dlDir}, evB) // 接收端 + qA.eng, qB.eng = engA, engB + + done := make(chan struct{}) + defer close(done) + qA.run(t, done) + qB.run(t, done) + + const sid = "sess-1" + if err := engB.StartIncoming(sid, "A", ""); err != nil { // 先武装接收端 + t.Fatalf("StartIncoming: %v", err) + } + if err := engA.StartOutgoing(sid, "B", srcPath, ""); err != nil { + t.Fatalf("StartOutgoing: %v", err) + } + + var savedPath string + select { + case savedPath = <-evB.saved: + case <-evB.failed: + t.Fatal("receiver reported failed") + case <-evA.failed: + t.Fatal("sender reported failed") + case <-time.After(20 * time.Second): + t.Fatal("timeout waiting for transfer to complete") + } + + got, err := os.ReadFile(savedPath) + if err != nil { + t.Fatalf("read saved file: %v", err) + } + if len(got) != len(want) { + t.Fatalf("size mismatch: got %d want %d", len(got), len(want)) + } + if !bytes.Equal(got, want) { + t.Fatal("content mismatch between source and received file") + } + + // 给发送端 waitAck/收尾一点时间收束,避免泄漏 goroutine 噪声(非断言项)。 + time.Sleep(200 * time.Millisecond) +} + +// TestParseICEServers 验证 urls 单串 / 数组两种形态都能解析(与 web getICEServers 输出兼容)。 +func TestParseICEServers(t *testing.T) { + out := parseICEServers(`[{"urls":"stun:a:1"},{"urls":["turns:b:2","turn:b:3"],"username":"u","credential":"c"}]`) + if len(out) != 2 { + t.Fatalf("want 2 servers, got %d", len(out)) + } + if len(out[0].URLs) != 1 || out[0].URLs[0] != "stun:a:1" { + t.Errorf("server0 urls wrong: %+v", out[0].URLs) + } + if len(out[1].URLs) != 2 || out[1].Username != "u" { + t.Errorf("server1 wrong: %+v", out[1]) + } +} diff --git a/desktop/engine/session.go b/desktop/engine/session.go new file mode 100644 index 0000000..859ef29 --- /dev/null +++ b/desktop/engine/session.go @@ -0,0 +1,524 @@ +package engine + +import ( + "encoding/json" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "sync" + "sync/atomic" + "time" + + "github.com/pion/ice/v4" + "github.com/pion/webrtc/v4" +) + +const ( + roleSender = "sender" + roleReceiver = "receiver" +) + +// session 是一条 P2P 传输的全部状态:pion 连接 + 数据通道 + 角色相关的收发循环。 +type session struct { + eng *Engine + id string + peer string + role string + + pc *webrtc.PeerConnection + dc *webrtc.DataChannel + + // 发送端 + filePath string + fileSize int64 + + // 计数器(跨 goroutine:发送循环 / pion 回调)。 + bytesSent atomic.Int64 // 已塞进本地 SCTP 缓冲的字节 + received atomic.Int64 // 接收端已写盘字节 + acked atomic.Int64 // 接收端经 ack 回传的已收字节(发送端进度/完成依据) + + // 接收端落盘 + out *os.File + outName string + + bufLow chan struct{} // OnBufferedAmountLow 唤醒发送循环 + done chan struct{} // 会话终结信号 + closeOnce sync.Once + + stateMu sync.Mutex + state string +} + +func newSession(e *Engine, id, peer, role string, pc *webrtc.PeerConnection) *session { + return &session{ + eng: e, + id: id, + peer: peer, + role: role, + pc: pc, + bufLow: make(chan struct{}, 1), + done: make(chan struct{}), + } +} + +// newPeerConnection 建一个配了大 SCTP 接收窗的 pion 连接。SetSCTPMaxReceiveBufferSize 把 +// 接收窗顶到 4MB(远高于 WebKit 写死的 256KB),是吞吐专项的核心——使本端作接收端时 +// `吞吐 ≈ rwnd/RTT` 的天花板抬高一个量级。pion 默认用真实 IP 的 host 候选(不做 mDNS +// 混淆),故同内网天然落到 host↔host 直连,避开 mac WKWebView 的 srflx 退化。 +func newPeerConnection(iceServersJSON string) (*webrtc.PeerConnection, error) { + se := webrtc.SettingEngine{} + se.SetSCTPMaxReceiveBufferSize(sctpReceiveBuffer) + // 解析对端的 mDNS(.local)候选。iOS / Safari 出于隐私只播 mDNS host 候选;pion 默认不 + // 解析,便无法与之成 host↔host、退而走 srflx/relay 慢路径(实测 mac→iOS 落到 relay↔srflx、 + // 1MB/s 且停滞)。QueryOnly=解析对端 mDNS、本端仍播真实 IP host 候选(对端可直接用), + // 双向打通同内网直连——这是吞吐的前提(host↔host 低 RTT 下 iOS 的 256KB rwnd 才非约束)。 + se.SetICEMulticastDNSMode(ice.MulticastDNSModeQueryOnly) + api := webrtc.NewAPI(webrtc.WithSettingEngine(se)) + pc, err := api.NewPeerConnection(webrtc.Configuration{ICEServers: parseICEServers(iceServersJSON)}) + if err != nil { + return nil, fmt.Errorf("new peerconnection: %w", err) + } + return pc, nil +} + +// wireCommon 绑两端共用的回调:trickle ICE 候选外发、连接状态变化。 +func (s *session) wireCommon() { + s.pc.OnICECandidate(func(c *webrtc.ICECandidate) { + if c == nil { + s.eng.ev.OnLog(fmt.Sprintf("session %s local gathering done", s.id)) + return + } + s.eng.ev.OnLog(fmt.Sprintf("session %s cand local: %s/%s", s.id, c.Typ, c.Protocol)) + init := c.ToJSON() + s.emitSignal(s.peer, nil, &init) + }) + // ICE 连接态变化日志:诊断「卡在等待对方接受 → 退中继」时连到哪一步(checking/connected/ + // failed/disconnected)。 + s.pc.OnICEConnectionStateChange(func(st webrtc.ICEConnectionState) { + s.eng.ev.OnLog(fmt.Sprintf("session %s ice: %s", s.id, st)) + }) + s.pc.OnConnectionStateChange(func(st webrtc.PeerConnectionState) { + switch st { + case webrtc.PeerConnectionStateConnected: + s.setState("connected") + s.logSelectedPair() + case webrtc.PeerConnectionStateFailed: + s.fail(errors.New("peerconnection failed")) + case webrtc.PeerConnectionStateClosed: + s.setState("closed") + } + }) +} + +// --- 发送端 --- + +func (s *session) bindSenderDC(dc *webrtc.DataChannel) { + s.dc = dc + dc.SetBufferedAmountLowThreshold(lowWatermark) + dc.OnBufferedAmountLow(func() { + select { + case s.bufLow <- struct{}{}: + default: + } + }) + dc.OnMessage(func(msg webrtc.DataChannelMessage) { + if !msg.IsString { + return + } + var a ackFrame + if json.Unmarshal(msg.Data, &a) == nil && a.Type == "ack" { + if a.Bytes > s.acked.Load() { + s.acked.Store(a.Bytes) + s.emitProgress() + } + } + }) + dc.OnOpen(func() { go s.streamFile() }) +} + +// streamFile 把文件按 64KB 分片经数据通道发出,水位回压与 web 引擎一致:缓冲将越 HIGH 前 +// 等抽干到 LOW。发完 done 帧后先抽干本地缓冲、再等接收端 ack 追平总量才宣告 completed +// (对齐接收端真实收程,避免发送端抢先完成)。 +func (s *session) streamFile() { + f, err := os.Open(s.filePath) + if err != nil { + s.fail(fmt.Errorf("open %s: %w", s.filePath, err)) + return + } + defer f.Close() + + meta, _ := json.Marshal(metaFrame{Type: "meta", Name: filepath.Base(s.filePath), Size: s.fileSize}) + if err := s.dc.SendText(string(meta)); err != nil { + s.fail(fmt.Errorf("send meta: %w", err)) + return + } + + buf := make([]byte, chunkSize) + for { + select { + case <-s.done: + return + default: + } + n, rerr := f.Read(buf) + if n > 0 { + if s.dc.BufferedAmount()+uint64(n) > highWatermark { + s.waitBufferLow(lowWatermark) + } + if err := s.dc.Send(buf[:n]); err != nil { + s.fail(fmt.Errorf("send chunk: %w", err)) + return + } + s.bytesSent.Add(int64(n)) + s.emitProgress() + } + if rerr == io.EOF { + break + } + if rerr != nil { + s.fail(fmt.Errorf("read %s: %w", s.filePath, rerr)) + return + } + } + + doneB, _ := json.Marshal(doneFrame{Type: "done"}) + _ = s.dc.SendText(string(doneB)) + s.waitDrained() + s.waitAck(s.fileSize) + s.setState("completed") +} + +// --- 接收端 --- + +func (s *session) bindReceiverDC(dc *webrtc.DataChannel) { + s.dc = dc + dc.OnMessage(func(msg webrtc.DataChannelMessage) { + if msg.IsString { + s.handleControl(msg.Data) + return + } + s.writeChunk(msg.Data) + }) + go s.ackLoop() +} + +func (s *session) handleControl(data []byte) { + var probe struct { + Type string `json:"type"` + } + if json.Unmarshal(data, &probe) != nil { + return + } + switch probe.Type { + case "meta": + var m metaFrame + if json.Unmarshal(data, &m) == nil { + s.openOutput(m.Name) + } + case "done": + s.finalize() + } +} + +func (s *session) openOutput(name string) { + s.outName = name + tmp := filepath.Join(s.eng.downloadDir(), "."+s.id+".part") + f, err := os.Create(tmp) + if err != nil { + s.fail(fmt.Errorf("create %s: %w", tmp, err)) + return + } + s.out = f +} + +// writeChunk 直接写盘——无 base64 桥、无 OPFS、无 WebView 内存。pion 串行投递同一通道的 +// OnMessage,故写入天然有序。 +func (s *session) writeChunk(b []byte) { + if s.out != nil { + if _, err := s.out.Write(b); err != nil { + s.fail(fmt.Errorf("write chunk: %w", err)) + return + } + } + s.received.Add(int64(len(b))) + s.emitProgress() +} + +// ackLoop 每 200ms 把最新已收字节回传发送端(节流,避免每片都回与下行争用通道)。 +func (s *session) ackLoop() { + t := time.NewTicker(ackIntervalMs * time.Millisecond) + defer t.Stop() + var lastSent int64 = -1 + for { + select { + case <-s.done: + return + case <-t.C: + if r := s.received.Load(); r != lastSent { + s.sendAck(r) + lastSent = r + } + } + } +} + +func (s *session) sendAck(bytes int64) { + if s.dc == nil || s.dc.ReadyState() != webrtc.DataChannelStateOpen { + return + } + b, _ := json.Marshal(ackFrame{Type: "ack", Bytes: bytes}) + _ = s.dc.SendText(string(b)) +} + +// finalize 在收到 done 帧时调用:done 意味所有分片已到(ordered),received 已等于总量。 +// 立即回最终 ack、关闭文件、改名落到不冲突路径,报告路径与 completed。 +func (s *session) finalize() { + s.sendAck(s.received.Load()) + if s.out != nil { + s.out.Close() + dir := s.eng.downloadDir() + final := nonCollidingPath(dir, s.outName) + tmp := filepath.Join(dir, "."+s.id+".part") + if err := os.Rename(tmp, final); err != nil { + s.fail(fmt.Errorf("rename to %s: %w", final, err)) + return + } + s.out = nil + s.eng.ev.OnSaved(s.id, final) + } + s.setState("completed") +} + +// --- 信令 --- + +func (s *session) handleSignal(data []byte) error { + var p signalPayload + if err := json.Unmarshal(data, &p); err != nil { + return fmt.Errorf("decode signal: %w", err) + } + s.eng.ev.OnLog(fmt.Sprintf("session %s signal in: %s", s.id, p.Type)) + switch p.Type { + case "offer": + if p.SDP == nil { + return errors.New("offer missing sdp") + } + if err := s.pc.SetRemoteDescription(*p.SDP); err != nil { + return fmt.Errorf("set remote (offer): %w", err) + } + ans, err := s.pc.CreateAnswer(nil) + if err != nil { + return fmt.Errorf("create answer: %w", err) + } + if err := s.pc.SetLocalDescription(ans); err != nil { + return fmt.Errorf("set local (answer): %w", err) + } + s.emitSignal(s.peer, &ans, nil) + case "answer": + if p.SDP == nil { + return errors.New("answer missing sdp") + } + if err := s.pc.SetRemoteDescription(*p.SDP); err != nil { + return fmt.Errorf("set remote (answer): %w", err) + } + case "ice": + if p.Candidate == nil { + return nil + } + if err := s.pc.AddICECandidate(*p.Candidate); err != nil { + return fmt.Errorf("add ice candidate: %w", err) + } + } + return nil +} + +// emitSignal 把一条出站信令交给宿主 POST。type 由载荷推导:有 candidate 即 ice,否则取 +// sdp 自身的类型(offer/answer),与 p2p.ts 的 payload.type 语义一致。 +func (s *session) emitSignal(toPeer string, sdp *webrtc.SessionDescription, cand *webrtc.ICECandidateInit) { + p := signalPayload{SDP: sdp, Candidate: cand} + switch { + case cand != nil: + p.Type = "ice" + case sdp != nil: + p.Type = sdp.Type.String() + } + b, err := json.Marshal(p) + if err != nil { + return + } + s.eng.ev.OnSignal(s.id, toPeer, string(b)) +} + +// --- 进度 / 状态 / 收尾 --- + +func (s *session) emitProgress() { + var b int64 + if s.role == roleSender { + if a := s.acked.Load(); a > 0 { + b = min(s.fileSize, a) + } else { + b = s.bytesSent.Load() - int64(s.dc.BufferedAmount()) + if b < 0 { + b = 0 + } + } + } else { + b = s.received.Load() + } + s.eng.ev.OnProgress(s.id, b) +} + +// waitBufferLow 等本地 SCTP 缓冲落到 threshold 以下:靠 OnBufferedAmountLow 事件唤醒, +// 兼一个 50ms 轮询兜底(事件阈值设在 lowWatermark,其它 threshold 靠轮询)。 +func (s *session) waitBufferLow(threshold uint64) { + start := time.Now() + warned := false + for { + if s.dc.BufferedAmount() <= threshold { + return + } + select { + case <-s.bufLow: + case <-time.After(50 * time.Millisecond): + case <-s.done: + return + } + // 停滞看门狗:本地缓冲长时间不抽干=SCTP 不再向网络交付,多为接收端 rwnd=0 + // (主线程/落盘卡住)或中继路径拥塞窗口塌缩。打一条诊断日志(不中断,polling 接力)。 + if !warned && time.Since(start) > 5*time.Second { + warned = true + s.eng.ev.OnLog(fmt.Sprintf( + "session %s send stalled >5s: buffered=%d sent=%d acked=%d", + s.id, s.dc.BufferedAmount(), s.bytesSent.Load(), s.acked.Load())) + } + } +} + +// logSelectedPair 在连接建立后异步取 pion 选中的候选对类型(host/srflx/prflx/relay)并打日志, +// 用于确认是否走上同内网 host↔host 直连。nominated pair 在 connected 时常未定,轮询几次。 +func (s *session) logSelectedPair() { + go func() { + for i := 0; i < 12; i += 1 { + if local, remote, ok := s.selectedPair(); ok { + s.eng.ev.OnLog(fmt.Sprintf("session %s ice pair: local=%s remote=%s", s.id, local, remote)) + s.eng.ev.OnIcePair(s.id, local, remote) + return + } + select { + case <-time.After(300 * time.Millisecond): + case <-s.done: + return + } + } + }() +} + +func (s *session) selectedPair() (local, remote string, ok bool) { + defer func() { _ = recover() }() // 早期连接链路可能尚 nil,吞 panic + sctp := s.pc.SCTP() + if sctp == nil { + return "", "", false + } + dtls := sctp.Transport() + if dtls == nil { + return "", "", false + } + it := dtls.ICETransport() + if it == nil { + return "", "", false + } + pair, err := it.GetSelectedCandidatePair() + if err != nil || pair == nil || pair.Local == nil || pair.Remote == nil { + return "", "", false + } + local = fmt.Sprintf("%s/%s", pair.Local.Typ, pair.Local.Protocol) + remote = fmt.Sprintf("%s/%s", pair.Remote.Typ, pair.Remote.Protocol) + return local, remote, true +} + +// waitDrained 等本地缓冲彻底抽干(done 帧也已离开本端 SCTP)。 +func (s *session) waitDrained() { + for { + if s.dc.BufferedAmount() == 0 { + return + } + select { + case <-time.After(20 * time.Millisecond): + case <-s.done: + return + } + } +} + +// waitAck 等接收端 ack 追平 total(确认已全收)。超时兜底;无任何 ack(不应发生,因对端 +// 皆回 ack)则不空等。 +func (s *session) waitAck(total int64) { + if s.acked.Load() == 0 { + return + } + deadline := time.Now().Add(ackCompleteTimeoutSec * time.Second) + for s.acked.Load() < total { + if time.Now().After(deadline) { + return + } + select { + case <-time.After(50 * time.Millisecond): + case <-s.done: + return + } + } +} + +func (s *session) setState(st string) { + s.stateMu.Lock() + // 终态不可被覆盖:completed/failed 之后即便连接随后 closed/failed 也不回退——否则接收端 + // 完成后、发送端关连接时本端会把已完成会话误翻成 closed/failed。 + if s.state == st || s.state == "completed" || s.state == "failed" { + s.stateMu.Unlock() + return + } + s.state = st + s.stateMu.Unlock() + + s.eng.ev.OnState(s.id, st) + switch st { + case "failed", "closed": + s.cleanup() + case "completed": + // 不立即拆连接。接收端在 finalize 里刚发出最终 ack;若此刻就关通道,发送端(尤其在 + // 前向数据把链路打满、反向 ack 被饿的高吞吐场景)可能收不到最终 ack → 误判失败(实测 + // iOS→mac 实际成功却被 iOS 标失败、进度停在 210/271)。改为留连接:ackLoop 每 200ms 续发 + // 最终 ack,待前向洪流结束、反向腾出后送达发送端。真正拆连接交给 SSE 终态驱动的 Cancel + // (hub→p2pCleanup→P2PCancel);60s 兜底防泄漏。 + go func() { + select { + case <-time.After(60 * time.Second): + s.cleanup() + case <-s.done: + } + }() + } +} + +func (s *session) fail(err error) { + s.eng.ev.OnLog(fmt.Sprintf("session %s fail: %v", s.id, err)) + s.setState("failed") +} + +func (s *session) cleanup() { + s.closeOnce.Do(func() { + close(s.done) + if s.out != nil { + s.out.Close() + // 半截 .part 留作诊断证据由上层决定清理;此处仅关句柄。 + } + if s.dc != nil { + _ = s.dc.Close() + } + if s.pc != nil { + _ = s.pc.Close() + } + s.eng.remove(s.peer, s) + }) +} diff --git a/desktop/engine/wire.go b/desktop/engine/wire.go new file mode 100644 index 0000000..984bdb3 --- /dev/null +++ b/desktop/engine/wire.go @@ -0,0 +1,120 @@ +// Package engine 是 cdrop 的原生数据面:用 pion/webrtc 在 Go 进程内完成 P2P 文件收发, +// 取代 WebView 内 JS 引擎跑 WebRTC 的旧路径。动机见 desktop/NATIVE-TRANSFER.md——逃离 +// WebKit 写死的 256KB SCTP 接收窗 + 渲染器节流,并让桌面与 iOS(经 gomobile)共用同一份实现。 +// +// 本包刻意只依赖 pion + 标准库,不碰 Wails / 平台代码:① 将来 `gomobile bind` 出 iOS +// xcframework 时不会拖入不可编译的桌面依赖;② 一切 HTTP(信令 / 状态机 POST / ICE 凭据) +// 留给宿主(桌面 JS / iOS Swift),引擎只经回调收发不透明信令串——纯数据面。 +// +// 线协议与 web 引擎 web/src/features/transfer/p2p.ts 逐字节一致,故 Go 端可与浏览器 / +// iOS 的 JS 引擎互通:DataChannel "cdrop-file"(ordered),控制帧 meta/done/ack 走文本帧、 +// 文件分片走二进制帧;信令 payload 形如 {type, sdp?, candidate?}。 +package engine + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + + "github.com/pion/webrtc/v4" +) + +const ( + channelName = "cdrop-file" + chunkSize = 64 * 1024 + highWatermark = 16 * 1024 * 1024 + lowWatermark = 4 * 1024 * 1024 + // 接收端把已收字节回传发送端的节流间隔(毫秒),与 web 引擎 ACK_INTERVAL_MS 一致。 + ackIntervalMs = 200 + // 发送端等接收端 ack 追平总量的上限(秒):超时仍照常收尾,避免 ack 丢失致悬挂。 + ackCompleteTimeoutSec = 30 + // pion 接收窗口:远高于 WebKit 写死的 256KB,使 pion 作接收端不被 rwnd × RTT 掐死 + // (吞吐专项的核心修复,见 SetSCTPMaxReceiveBufferSize)。 + sctpReceiveBuffer = 4 * 1024 * 1024 +) + +// 控制帧(文本帧),与 p2p.ts 的 meta/done/ack 同形。发送端 meta 仅含 name/size(sha256 +// 在 streamFile 路径上 web 端亦不发),故此处不带。 +type metaFrame struct { + Type string `json:"type"` + Name string `json:"name"` + Size int64 `json:"size"` +} + +type doneFrame struct { + Type string `json:"type"` +} + +type ackFrame struct { + Type string `json:"type"` + Bytes int64 `json:"bytes"` +} + +// signalPayload 与 p2p.ts 的 SignalPayload 同形:offer/answer 带 sdp(RTCSessionDescriptionInit), +// ice 带 candidate(RTCIceCandidateInit)。pion 的 SessionDescription / ICECandidateInit 的 JSON +// 标签与浏览器侧一致,故双向可直接 (un)marshal。 +type signalPayload struct { + Type string `json:"type"` + SDP *webrtc.SessionDescription `json:"sdp,omitempty"` + Candidate *webrtc.ICECandidateInit `json:"candidate,omitempty"` +} + +// iceServerJSON 解析 web getICEServers() 的输出:urls 既可能是单串也可能是数组。 +type iceServerJSON struct { + URLs json.RawMessage `json:"urls"` + Username string `json:"username,omitempty"` + Credential string `json:"credential,omitempty"` +} + +// parseICEServers 把宿主传来的 ICE 服务器 JSON(与 web 同源,来自 /api/calls/credentials) +// 解析成 pion 的配置。urls 兼容字符串与字符串数组两种形态。 +func parseICEServers(raw string) []webrtc.ICEServer { + if strings.TrimSpace(raw) == "" { + return nil + } + var arr []iceServerJSON + if err := json.Unmarshal([]byte(raw), &arr); err != nil { + return nil + } + out := make([]webrtc.ICEServer, 0, len(arr)) + for _, s := range arr { + var urls []string + var one string + if json.Unmarshal(s.URLs, &one) == nil { + urls = []string{one} + } else { + _ = json.Unmarshal(s.URLs, &urls) + } + if len(urls) == 0 { + continue + } + srv := webrtc.ICEServer{URLs: urls} + if s.Username != "" { + srv.Username = s.Username + } + if s.Credential != "" { + srv.Credential = s.Credential + } + out = append(out, srv) + } + return out +} + +// nonCollidingPath 返回 dir 下不与现有文件冲突的目标路径:name、name (1)、name (2)…… +// 接收端落盘时用,避免覆盖既有文件(对齐 SaveDownload 的不覆盖语义)。 +func nonCollidingPath(dir, name string) string { + p := filepath.Join(dir, name) + if _, err := os.Stat(p); os.IsNotExist(err) { + return p + } + ext := filepath.Ext(name) + base := strings.TrimSuffix(name, ext) + for i := 1; ; i += 1 { + c := filepath.Join(dir, fmt.Sprintf("%s (%d)%s", base, i, ext)) + if _, err := os.Stat(c); os.IsNotExist(err) { + return c + } + } +} diff --git a/desktop/go.mod b/desktop/go.mod index f38fa6e..2db033a 100644 --- a/desktop/go.mod +++ b/desktop/go.mod @@ -1,15 +1,17 @@ module cdrop-desktop -go 1.24 +go 1.24.0 require ( fyne.io/systray v1.12.2 git.sr.ht/~jackmordaunt/go-toast/v2 v2.0.3 github.com/ebitengine/purego v0.10.1 + github.com/pion/ice/v4 v4.2.7 + github.com/pion/webrtc/v4 v4.2.15 github.com/wailsapp/wails/v2 v2.12.0 github.com/zalando/go-keyring v0.2.8 golang.design/x/clipboard v0.8.0 - golang.org/x/sys v0.33.0 + golang.org/x/sys v0.41.0 ) require ( @@ -28,6 +30,20 @@ require ( github.com/leaanthony/u v1.1.1 // indirect github.com/mattn/go-colorable v0.1.13 // indirect github.com/mattn/go-isatty v0.0.20 // indirect + github.com/pion/datachannel v1.6.0 // indirect + github.com/pion/dtls/v3 v3.1.4 // indirect + github.com/pion/interceptor v0.1.45 // indirect + github.com/pion/logging v0.2.4 // indirect + github.com/pion/mdns/v2 v2.1.0 // indirect + github.com/pion/randutil v0.1.0 // indirect + github.com/pion/rtcp v1.2.16 // indirect + github.com/pion/rtp v1.10.2 // indirect + github.com/pion/sctp v1.10.0 // indirect + github.com/pion/sdp/v3 v3.0.18 // indirect + github.com/pion/srtp/v3 v3.0.11 // indirect + github.com/pion/stun/v3 v3.1.5 // indirect + github.com/pion/transport/v4 v4.0.2 // indirect + github.com/pion/turn/v5 v5.0.9 // indirect github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c // indirect github.com/pkg/errors v0.9.1 // indirect github.com/rivo/uniseg v0.4.7 // indirect @@ -37,13 +53,15 @@ require ( github.com/valyala/fasttemplate v1.2.2 // indirect github.com/wailsapp/go-webview2 v1.0.22 // indirect github.com/wailsapp/mimetype v1.4.1 // indirect + github.com/wlynxg/anet v0.0.5 // indirect golang.design/x/x11 v0.2.0 // indirect - golang.org/x/crypto v0.33.0 // indirect + golang.org/x/crypto v0.48.0 // indirect golang.org/x/exp/shiny v0.0.0-20250606033433-dcc06ee1d476 // indirect golang.org/x/image v0.28.0 // indirect golang.org/x/mobile v0.0.0-20250606033058-a2a15c67f36f // indirect - golang.org/x/net v0.35.0 // indirect - golang.org/x/text v0.26.0 // indirect + golang.org/x/net v0.50.0 // indirect + golang.org/x/text v0.34.0 // indirect + golang.org/x/time v0.14.0 // indirect ) // replace github.com/wailsapp/wails/v2 v2.12.0 => /Users/commilitia/go/pkg/mod diff --git a/desktop/go.sum b/desktop/go.sum index df7a571..aa94936 100644 --- a/desktop/go.sum +++ b/desktop/go.sum @@ -42,6 +42,40 @@ github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovk github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/pion/datachannel v1.6.0 h1:XecBlj+cvsxhAMZWFfFcPyUaDZtd7IJvrXqlXD/53i0= +github.com/pion/datachannel v1.6.0/go.mod h1:ur+wzYF8mWdC+Mkis5Thosk+u/VOL287apDNEbFpsIk= +github.com/pion/dtls/v3 v3.1.4 h1:QhvtMflMfu9Kf0RcDC5BJBle4caPskByrKQR6uuYqpY= +github.com/pion/dtls/v3 v3.1.4/go.mod h1:cr/qotLISUw/9C1m83ZPNZtj9WnXkYLpfCptPqbkInc= +github.com/pion/ice/v4 v4.2.7 h1:zDEbC6MiEdhQpF8TxBOTws+NU6ZgGpveHrQq4Lc1kao= +github.com/pion/ice/v4 v4.2.7/go.mod h1:9SNPaq0c7El/ki8leJzyCkK10zsskprR3zTNbO3monY= +github.com/pion/interceptor v0.1.45 h1:6PUo/5829bIfRFIPPJQzuDn8EjxRTSB/CSD7QVCOaqo= +github.com/pion/interceptor v0.1.45/go.mod h1:gNDYM/uFKcLe/B3gS2/7+aw6z+RDiMy2qKTnF1LO31w= +github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8= +github.com/pion/logging v0.2.4/go.mod h1:DffhXTKYdNZU+KtJ5pyQDjvOAh/GsNSyv1lbkFbe3so= +github.com/pion/mdns/v2 v2.1.0 h1:3IJ9+Xio6tWYjhN6WwuY142P/1jA0D5ERaIqawg/fOY= +github.com/pion/mdns/v2 v2.1.0/go.mod h1:pcez23GdynwcfRU1977qKU0mDxSeucttSHbCSfFOd9A= +github.com/pion/randutil v0.1.0 h1:CFG1UdESneORglEsnimhUjf33Rwjubwj6xfiOXBa3mA= +github.com/pion/randutil v0.1.0/go.mod h1:XcJrSMMbbMRhASFVOlj/5hQial/Y8oH/HVo7TBZq+j8= +github.com/pion/rtcp v1.2.16 h1:fk1B1dNW4hsI78XUCljZJlC4kZOPk67mNRuQ0fcEkSo= +github.com/pion/rtcp v1.2.16/go.mod h1:/as7VKfYbs5NIb4h6muQ35kQF/J0ZVNz2Z3xKoCBYOo= +github.com/pion/rtp v1.10.2 h1:l+f6tTDcAH6xwepaAoW791ddhuYsJlqRATOzirO04Mo= +github.com/pion/rtp v1.10.2/go.mod h1:Au8fc6cEByy8RLTwKTQTEeQqDB/SJDxwL4mZuxYA5Pk= +github.com/pion/sctp v1.10.0 h1:qeoD6swF/2M5bYRcAGayqSbTKX3m4AW29CiQxG1+Pfg= +github.com/pion/sctp v1.10.0/go.mod h1:N20Dq6LY+JvJDAh9VVh1JELngb2rQ8dPgds5yBWiPgw= +github.com/pion/sdp/v3 v3.0.18 h1:l0bAXazKHpepazVdp+tPYnrsy9dfh7ZbT8DxesH5ZnI= +github.com/pion/sdp/v3 v3.0.18/go.mod h1:ZREGo6A9ZygQ9XkqAj5xYCQtQpif0i6Pa81HOiAdqQ8= +github.com/pion/srtp/v3 v3.0.11 h1:GiESUr54/K4UuPigfq/CvWUed80JenQAHXn0C2MQQIQ= +github.com/pion/srtp/v3 v3.0.11/go.mod h1:EeZOi/sd6glM1EXapg051gdNWO9yWT1YSsgQ4SlJkns= +github.com/pion/stun/v3 v3.1.5 h1:Y1FHlhaI6+4UoC5i/zQf4F7JvdZtB24/05oyy/GF1x8= +github.com/pion/stun/v3 v3.1.5/go.mod h1:zRUghXSQU32Lx5orJsz3uYMkIihweXb3mu5gIns02fs= +github.com/pion/transport/v3 v3.1.1 h1:Tr684+fnnKlhPceU+ICdrw6KKkTms+5qHMgw6bIkYOM= +github.com/pion/transport/v3 v3.1.1/go.mod h1:+c2eewC5WJQHiAA46fkMMzoYZSuGzA/7E2FPrOYHctQ= +github.com/pion/transport/v4 v4.0.2 h1:ifYlPqNwsy6aKQ9y8yzxXlHae5431ZrH2avkD/Rn6Tk= +github.com/pion/transport/v4 v4.0.2/go.mod h1:06hFI+jCFcok2X2MekVufNZ/uzNZXivGBPfviSVcjgM= +github.com/pion/turn/v5 v5.0.9 h1:zNeBfRyzGn7MPyUTvmvxeltLEjlFdSLPT1tlakoaOXM= +github.com/pion/turn/v5 v5.0.9/go.mod h1:u3XjBqy2Z4+NhCUpDoOSsNuQDrPLvKStlCGWk6sTQ1E= +github.com/pion/webrtc/v4 v4.2.15 h1:Ir/MauNFCfg+kgyBYPQLiGdVWFlzEcLxqtuzAkYkky0= +github.com/pion/webrtc/v4 v4.2.15/go.mod h1:CPTcyLfIzC4scOkQ4UY4pj6WvbUGhcNLIpK28cP5h6M= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c h1:+mdjkGKdHQG3305AYmdv1U2eRNDiU2ErMBj1gwrq8eQ= github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c/go.mod h1:7rwL4CYBLnjLxUqIJNnCWiEdr3bn6IUYi15bNlnbCCU= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -69,14 +103,16 @@ github.com/wailsapp/mimetype v1.4.1 h1:pQN9ycO7uo4vsUUuPeHEYoUkLVkaRntMnHJxVwYhw github.com/wailsapp/mimetype v1.4.1/go.mod h1:9aV5k31bBOv5z6u+QP8TltzvNGJPmNJD4XlAL3U+j3o= github.com/wailsapp/wails/v2 v2.12.0 h1:BHO/kLNWFHYjCzucxbzAYZWUjub1Tvb4cSguQozHn5c= github.com/wailsapp/wails/v2 v2.12.0/go.mod h1:mo1bzK1DEJrobt7YrBjgxvb5Sihb1mhAY09hppbibQg= +github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= +github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= github.com/zalando/go-keyring v0.2.8 h1:6sD/Ucpl7jNq10rM2pgqTs0sZ9V3qMrqfIIy5YPccHs= github.com/zalando/go-keyring v0.2.8/go.mod h1:tsMo+VpRq5NGyKfxoBVjCuMrG47yj8cmakZDO5QGii0= golang.design/x/clipboard v0.8.0 h1:6VEcH28wwcSgKc+vnxHHDWiRjrakTQIAJnPUrt3aOgg= golang.design/x/clipboard v0.8.0/go.mod h1:s0pwrtA3Q9fgnVtGDmP5ZK/pp55cQKB23esKsjwWhWM= golang.design/x/x11 v0.2.0 h1:Uiwu2guGihsJX/ZCzpoDPFz5gR/Qntm08mvoBCmRydo= golang.design/x/x11 v0.2.0/go.mod h1:/5q1mFkdc1rL8mvB7DsQFi6as4tIkBv4FXjcP07mrkE= -golang.org/x/crypto v0.33.0 h1:IOBPskki6Lysi0lo9qQvbxiQ+FvsCC/YWOecCHAixus= -golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M= +golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= +golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= golang.org/x/exp/shiny v0.0.0-20250606033433-dcc06ee1d476 h1:Wdx0vgH5Wgsw+lF//LJKmWOJBLWX6nprsMqnf99rYDE= golang.org/x/exp/shiny v0.0.0-20250606033433-dcc06ee1d476/go.mod h1:ygj7T6vSGhhm/9yTpOQQNvuAUFziTH7RUiH74EoE2C8= golang.org/x/image v0.28.0 h1:gdem5JW1OLS4FbkWgLO+7ZeFzYtL3xClb97GaUzYMFE= @@ -84,20 +120,22 @@ golang.org/x/image v0.28.0/go.mod h1:GUJYXtnGKEUgggyzh+Vxt+AviiCcyiwpsl8iQ8MvwGY golang.org/x/mobile v0.0.0-20250606033058-a2a15c67f36f h1:/n+PL2HlfqeSiDCuhdBbRNlGS/g2fM4OHufalHaTVG8= golang.org/x/mobile v0.0.0-20250606033058-a2a15c67f36f/go.mod h1:ESkJ836Z6LpG6mTVAhA48LpfW/8fNR0ifStlH2axyfg= golang.org/x/net v0.0.0-20210505024714-0287a6fb4125/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= -golang.org/x/net v0.35.0 h1:T5GQRQb2y08kTAByq9L4/bz8cipCdA8FbRTXewonqY8= -golang.org/x/net v0.35.0/go.mod h1:EglIi67kWsHKlRzzVMUD93VMSWGFOMSZgxFjparz1Qk= +golang.org/x/net v0.50.0 h1:ucWh9eiCGyDR3vtzso0WMQinm2Dnt8cFMuQa9K33J60= +golang.org/x/net v0.50.0/go.mod h1:UgoSli3F/pBgdJBHCTc+tp3gmrU4XswgGRgtnwWTfyM= golang.org/x/sys v0.0.0-20200810151505-1b9f1253b3ed/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw= -golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.26.0 h1:P42AVeLghgTYr4+xUnTRKDMqpar+PtX7KWuNQL21L8M= -golang.org/x/text v0.26.0/go.mod h1:QK15LZJUUQVJxhz7wXgxSy/CJaTFjd0G+YLonydOVQA= +golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= +golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= +golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI= +golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/desktop/platform/clipboard_darwin.m b/desktop/platform/clipboard_darwin.m index 5b2185f..3a4c7c6 100644 --- a/desktop/platform/clipboard_darwin.m +++ b/desktop/platform/clipboard_darwin.m @@ -14,6 +14,36 @@ long cdropPasteboardChangeCount(void) { // org.nspasteboard.{Concealed,Transient,AutoGenerated}Type — so the upload // policy can skip it. (Most password managers set no marker; this is a // best-effort secondary guard behind the server's short TTL.) +// cdropIsStagingPath 判断字符串是否为 Universal Clipboard / Handoff 的暂存文件路径——富文本 +// 复制时 pasteboard 的纯文本表示有时是 .../shared-pasteboard/.../xxx.rtfd 这样的路径而非内容, +// 绝不能当剪贴板文本上传同步。 +static BOOL cdropIsStagingPath(NSString *s) { + if (s == nil) { return NO; } + return ([s rangeOfString:@"/shared-pasteboard/"].location != NSNotFound) || + ([s rangeOfString:@"com.apple.coreservices.useractivityd"].location != NSNotFound); +} + +// cdropPlainFromRich 从 RTF / RTFD 富文本派生纯文本(复制富文本时纯文本类型可能缺失或为暂存 +// 路径)。取不到返回 nil。 +static NSString *cdropPlainFromRich(NSPasteboard *pb) { + NSData *data = [pb dataForType:NSPasteboardTypeRTFD]; + NSAttributedString *as = nil; + if (data) { + as = [[NSAttributedString alloc] initWithData:data + options:@{NSDocumentTypeDocumentAttribute: NSRTFDTextDocumentType} + documentAttributes:nil error:nil]; + } + if (as == nil) { + data = [pb dataForType:NSPasteboardTypeRTF]; + if (data) { + as = [[NSAttributedString alloc] initWithData:data + options:@{NSDocumentTypeDocumentAttribute: NSRTFTextDocumentType} + documentAttributes:nil error:nil]; + } + } + return as ? as.string : nil; +} + char *cdropPasteboardReadText(int *sensitive) { NSPasteboard *pb = [NSPasteboard generalPasteboard]; *sensitive = 0; @@ -26,7 +56,11 @@ char *cdropPasteboardReadText(int *sensitive) { } } NSString *s = [pb stringForType:NSPasteboardTypeString]; - if (s == nil) { return NULL; } + // 纯文本缺失或取到的是 Handoff 暂存路径 → 从 RTF/RTFD 派生真正的文本。 + if (s == nil || cdropIsStagingPath(s)) { + s = cdropPlainFromRich(pb); + } + if (s == nil || cdropIsStagingPath(s)) { return NULL; } // 仍拿不到文本:当非文本忽略 const char *utf8 = [s UTF8String]; if (utf8 == NULL) { return NULL; } return strdup(utf8); diff --git a/web/src/features/clipboard/clipboard.ts b/web/src/features/clipboard/clipboard.ts index c839a9f..6852490 100644 --- a/web/src/features/clipboard/clipboard.ts +++ b/web/src/features/clipboard/clipboard.ts @@ -81,8 +81,20 @@ export function utf8ByteLength(text: string): number // uploadClipboard 把一段纯文本作为剪贴板写入推到云端(PUT /api/clipboard)。 // 字节超限 / 空字符串 / 网络错误统一抛 Error,调用点显示给用户。 +// 富文本经 Universal Clipboard / Handoff 会被暂存为 Group Containers 下的 .rtfd 文件,原生侧 +// 读 pasteboard 字符串有时取到的是该暂存文件路径而非文本(实测:mac 复制富文本 → 同步出 +// 一条 .../shared-pasteboard/items/.../xxx.rtfd 路径)。这类暂存路径绝不能当剪贴板内容同步—— +// 在共享上传口拦掉,对三端(mac NSPasteboard / iOS UIPasteboard)一致生效。标记串极特异, +// 几乎不可能出现在用户真实复制的文本里。 +function isPasteboardStagingPath(s: string): boolean +{ + return s.includes("/shared-pasteboard/") + || s.includes("com.apple.coreservices.useractivityd"); +} + export async function uploadClipboard(content: string): Promise { + if (isPasteboardStagingPath(content)) { return; } if (utf8ByteLength(content) > CLIPBOARD_MAX_BYTES) { throw new Error(t("errors.clipboardOverflow", { max: CLIPBOARD_MAX_BYTES })); diff --git a/web/src/features/transfer/Composer.tsx b/web/src/features/transfer/Composer.tsx index 05d353d..ac1ba22 100644 --- a/web/src/features/transfer/Composer.tsx +++ b/web/src/features/transfer/Composer.tsx @@ -4,6 +4,8 @@ import { Button, DynText, Panel, Tabs, TextArea } from "../../ui/primitives"; import { toast } from "../../ui/feedback"; import { t } from "../../i18n"; import { formatBytes } from "../../utils/format"; +import { isDesktop, nativeFileSource, pickFileForSend } from "../../net/desktop"; +import { fileSource, type FileSource } from "./source"; import s from "./Composer.module.css"; export type ComposerMode = "file" | "message"; @@ -11,8 +13,8 @@ export type ComposerMode = "file" | "message"; export interface ComposerProps { targetDevice: string | null; - selectedFile: File | null; - onFileChange: (f: File | null) => void; + selectedFile: FileSource | null; + onFileChange: (f: FileSource | null) => void; onSendFile: () => Promise | void; onSendMessage: (text: string) => Promise | void; /** 默认 'file'。 */ @@ -148,8 +150,8 @@ export function Composer(props: ComposerProps) interface FilePaneProps { - selectedFile: File | null; - onFileChange: (f: File | null) => void; + selectedFile: FileSource | null; + onFileChange: (f: FileSource | null) => void; fileInputRef: React.RefObject; } @@ -170,8 +172,23 @@ function FilePane(props: FilePaneProps) { e.preventDefault(); setDragOver(false); + // 拖放给的是浏览器 File(桌面亦无绝对路径)→ 包成 FileSource。桌面此路无 path、 + // 回退 JS 数据面;走原生数据面需经下方 handlePick 的原生文件对话框取路径。 const f = e.dataTransfer.files[0]; - if (f) { onFileChange(f); } + if (f) { onFileChange(fileSource(f)); } + }; + + // 取文件:桌面走原生文件对话框(Go 返回绝对路径 → nativeFileSource,发送走原生数据面, + // 逃离 256KB rwnd 与渲染器节流);浏览器 / iOS 走隐藏的 。 + const handlePick = async () => + { + if (isDesktop()) + { + const picked = await pickFileForSend(t("home.file.dropzone")); + if (picked) { onFileChange(nativeFileSource(picked, "")); } + return; + } + fileInputRef.current?.click(); }; const cls = [ @@ -184,7 +201,7 @@ function FilePane(props: FilePaneProps) <>
fileInputRef.current?.click()} + onClick={() => void handlePick()} onDragOver={onDragOver} onDragLeave={() => setDragOver(false)} onDrop={onDrop} @@ -194,7 +211,7 @@ function FilePane(props: FilePaneProps) { if (e.key === "Enter" || e.key === " ") { e.preventDefault(); - fileInputRef.current?.click(); + void handlePick(); } }} > @@ -227,7 +244,7 @@ function FilePane(props: FilePaneProps) onChange={(e) => { const f = e.currentTarget.files?.[0]; - if (f) { onFileChange(f); } + if (f) { onFileChange(fileSource(f)); } e.currentTarget.value = ""; }} /> diff --git a/web/src/features/transfer/p2p.ts b/web/src/features/transfer/p2p.ts index 4e8b066..4980b28 100644 --- a/web/src/features/transfer/p2p.ts +++ b/web/src/features/transfer/p2p.ts @@ -1,7 +1,15 @@ import { apiFetch } from "../../net/api"; import { getICEServers } from "./iceServers"; +import { isDesktop } from "../../net/desktop"; import { useAppStore, type CandidateBreakdown, type IceStats, type TransferPhase } from "../../store"; import { deliverIncoming, openIncomingSink, type IncomingSink } from "./incomingSink"; +import { + nativeCleanup, + nativeHandleSignal, + nativeHasPeer, + nativeStartIncoming, + nativeStartOutgoing, +} from "./p2pNative"; import type { FileSource } from "./source"; // WebRTC client wired to the brief §2 invariants: @@ -889,6 +897,8 @@ const p2pSessions = new Map(); export function p2pHandleSignal(from: string, payload: SignalPayload): void { + // 桌面:该对端的会话在原生后端(Go/pion)→ 把入站信令转给它。 + if (nativeHasPeer(from)) { nativeHandleSignal(from, payload); return; } const sess = p2pSessions.get(from); if (!sess) { @@ -901,6 +911,7 @@ export function p2pHandleSignal(from: string, payload: SignalPayload): void export function p2pCleanup(peerName: string): void { + if (nativeHasPeer(peerName)) { nativeCleanup(peerName); return; } const sess = p2pSessions.get(peerName); if (sess) { sess.cancel(); } p2pSessions.delete(peerName); @@ -937,6 +948,12 @@ export function p2pStartOutgoing( src: FileSource, ): P2PSession { + // 桌面且源带绝对路径(经原生文件选择):走原生数据面(Go/pion 直接读盘发送),逃离 + // WebView 的 256KB rwnd 与渲染器节流。无 path 的源(浏览器 File 拖放)仍走 JS 引擎。 + if (isDesktop() && src.path) + { + return nativeStartOutgoing(sessionId, receiverName, src.path, src.size); + } const sess = new Session(sessionId, receiverName, "sender"); p2pSessions.set(receiverName, sess); sess.stateListeners.add((s) => updateStoreState(sessionId, s)); @@ -955,6 +972,11 @@ export function p2pStartOutgoing( export function p2pStartIncoming(sessionId: string, senderName: string): P2PSession { + // 桌面:接收端一律走原生数据面(Go 直接写盘,无 base64 桥 / 无 OPFS),逃离 256KB rwnd。 + if (isDesktop()) + { + return nativeStartIncoming(sessionId, senderName); + } const sess = new Session(sessionId, senderName, "receiver"); p2pSessions.set(senderName, sess); sess.stateListeners.add((s) => updateStoreState(sessionId, s)); diff --git a/web/src/features/transfer/p2pNative.ts b/web/src/features/transfer/p2pNative.ts new file mode 100644 index 0000000..006bff6 --- /dev/null +++ b/web/src/features/transfer/p2pNative.ts @@ -0,0 +1,320 @@ +// p2pNative —— 桌面原生数据面后端。实现与 p2p.ts 同形的 P2PSession 接口与四个入口, +// 把 WebRTC 收发委派给 Go/pion 引擎(经 Wails 桥),逃离 WebView 的 256KB rwnd 与渲染器 +// 节流(见 desktop/NATIVE-TRANSFER.md)。 +// +// 分工:Go 引擎只跑纯 WebRTC + 文件 I/O;一切 HTTP 在此发——出站信令 POST /api/hub/signal、 +// 状态机 POST(/p2p 于 connected·发送端、/done 于 completed·接收端、/fail 于 failed)。引擎 +// 经 p2p:* 事件把进度 / 状态 / 出站信令 / 落盘路径反向回来。 +// +// p2p.ts 在桌面按对端把这些会话登记到原生后端(见其四个入口的分支),故 transfer.ts / +// hub.ts 一字不改。 +import { apiFetch } from "../../net/api"; +import { t } from "../../i18n"; +import { + nativeP2PCancel, + nativeP2PHandleSignal, + nativeP2PStartIncoming, + nativeP2PStartOutgoing, + subscribeNativeP2P, +} from "../../net/desktop"; +import { useAppStore, type TransferPhase } from "../../store"; +import { toast } from "../../ui/feedback"; +import { getICEServers } from "./iceServers"; +import type { P2PProgressEvent, P2PSession, P2PState } from "./p2p"; + +interface NativeSess +{ + sessionId: string; + peerName: string; + role: "sender" | "receiver"; + state: P2PState; + bytes: number; + total: number; + stateListeners: Set<(s: P2PState) => void>; + progressListeners: Set<(e: P2PProgressEvent) => void>; + p2pPosted: boolean; + // 信令串行链:seed 为 StartOutgoing/Incoming 的桥调用(确保 Go 侧会话已建),其后每条 + // 入站信令串到链尾、按到达顺序投给 Go——pion 在 SetRemoteDescription 前 AddICECandidate + // 会报错,故 offer/answer 必须先于其后 trickle 的 ICE 候选。等价于 Go 端测试里的 sigQueue。 + chain: Promise; +} + +const sessionsById = new Map(); +const idByPeer = new Map(); +let subscribed = false; + +// ensureSubscribed 懒注册全局 p2p:* 事件监听(仅一次),据 sessionId 分发到各会话。 +function ensureSubscribed(): void +{ + if (subscribed) { return; } + subscribed = true; + subscribeNativeP2P({ + onProgress: (sid, bytes) => + { + const s = sessionsById.get(sid); + if (!s) { return; } + s.bytes = bytes; + const total = s.total || sessionFileSize(sid); + for (const cb of s.progressListeners) { cb({ bytes, total }); } + pushStoreBytes(sid, bytes); + }, + onState: (sid, state) => + { + const s = sessionsById.get(sid); + if (!s) { return; } + applyState(s, state as P2PState); + }, + onSignal: (_sid, toPeer, payloadJSON) => + { + // 各信令独立 POST、不串行——串行会让一条慢 POST 头阻塞其后所有 trickle 候选, + // 在弱网下致 ICE 无法在 30s 内连通 → 退中继(实测「等待对方接受」回退)。乱序到达 + // 对端只触发对端已 swallow 的 addIceCandidate 告警,无害。 + let payload: unknown; + try { payload = JSON.parse(payloadJSON); } + catch { return; } + void apiFetch("/api/hub/signal", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ to: toPeer, payload }), + }).then((r) => + { + if (!r.ok && r.status !== 410) + { + // eslint-disable-next-line no-console + console.warn("native p2p signal HTTP", r.status); + } + }).catch(() => { /* 信令尽力而为;丢失由 ICE 超时兜底 */ }); + }, + onSaved: (_sid, path) => + { + toast.ok(t("transfer.savedTo", { path })); + }, + onIcePair: (sid, local, remote) => + { + setIcePair(sid, local, remote); + }, + }); +} + +// setIcePair 把 pion 选中的候选对填进记录的 iceStats,使调试面板能展示「实际走哪条路径」 +// (原生后端此前不填,致进入传输后调试面板候选信息空白)。候选分类计数 pion 未细分,留空。 +function setIcePair(sessionId: string, local: string, remote: string): void +{ + const store = useAppStore.getState(); + const cur = store.activeTransfers[sessionId]; + if (!cur) { return; } + const empty = { host: 0, mdns: 0, srflx: 0, prflx: 0, relay: 0 }; + store.upsertTransfer({ + ...cur, + iceStats: { + gathering: "complete", + connection: "connected", + candidates: { local: { ...empty }, remote: { ...empty } }, + selectedPair: { local, remote }, + }, + }); +} + +function sessionFileSize(sessionId: string): number +{ + return useAppStore.getState().activeTransfers[sessionId]?.fileSize ?? 0; +} + +// applyState 把引擎状态落到监听者 + store,并发起对应的状态机 POST(HTTP 留 JS)。 +function applyState(s: NativeSess, state: P2PState): void +{ + if (s.state === state) { return; } + s.state = state; + for (const cb of s.stateListeners) { cb(state); } + updateStoreState(s.sessionId, state); + + if (state === "connected") + { + // 推进 phase(原生后端此前漏设,致 UI 状态停在初始 phase、且 bytes/速度文案被抑制)。 + setPhase(s.sessionId, "ice_connected"); + // 发送端在通道就绪时把传输标记为 P2P_ACTIVE,使随后接收端的 /done 合法。 + if (s.role === "sender" && !s.p2pPosted) + { + s.p2pPosted = true; + markP2PActive(s.sessionId); + } + } + else if (state === "completed") + { + // /done 由接收端权威发出;发送端只落本地终态(updateStoreState 已置 DONE)。 + if (s.role === "receiver") { markServerDone(s.sessionId, s.bytes); } + } + else if (state === "failed") + { + markServerFail(s.sessionId, "native_p2p_failed"); + } +} + +// updateStoreState 与 p2p.ts 同名函数等价(此处复制以免与 p2p.ts 形成运行时循环依赖)。 +function updateStoreState(sessionId: string, p2pState: P2PState): void +{ + const store = useAppStore.getState(); + const cur = store.activeTransfers[sessionId]; + if (!cur) { return; } + if (p2pState === "connected") { store.upsertTransfer({ ...cur, state: "P2P_ACTIVE" }); } + else if (p2pState === "completed") { store.completeTransfer(sessionId, "DONE"); } + else if (p2pState === "failed") { store.completeTransfer(sessionId, "FAILED"); } + // closed 不主动移 history:交由上层 / 中继流程决定终态(同 p2p.ts)。 +} + +function pushStoreBytes(sessionId: string, bytes: number): void +{ + const store = useAppStore.getState(); + const cur = store.activeTransfers[sessionId]; + if (!cur) { return; } + const grew = bytes > (cur.bytesTransferred ?? 0); + store.upsertTransfer({ + ...cur, + // 字节流动即「传输中」:FLOW_PHASES 据此放开 bytes/速度文案,状态文案也随之正确。 + phase: cur.phase === "completing" ? cur.phase : "transferring", + bytesTransferred: bytes, + lastProgressAt: grew ? Date.now() : cur.lastProgressAt, + }); +} + +// setPhase 更新记录的 phase(驱动 UI 状态文案 + bytes/速度文案的可见性);同值跳过。 +function setPhase(sessionId: string, phase: TransferPhase): void +{ + const store = useAppStore.getState(); + const cur = store.activeTransfers[sessionId]; + if (!cur || cur.phase === phase) { return; } + store.upsertTransfer({ ...cur, phase }); +} + +function markP2PActive(sessionId: string): void +{ + void apiFetch(`/api/transfer/${sessionId}/p2p`, { method: "POST" }).catch(() => { /* 409 无碍 */ }); +} + +function markServerDone(sessionId: string, bytes: number): void +{ + void apiFetch(`/api/transfer/${sessionId}/done`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ bytes_transferred: bytes }), + }).catch(() => { /* 接收端权威 /done;失败由 SSE 终态兜底 */ }); +} + +function markServerFail(sessionId: string, reason: string): void +{ + void apiFetch(`/api/transfer/${sessionId}/fail`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ reason }), + }).catch(() => { /* best-effort */ }); +} + +function makeSession(s: NativeSess): P2PSession +{ + return { + sessionId: s.sessionId, + peerName: s.peerName, + get state() { return s.state; }, + onState: (cb) => + { + s.stateListeners.add(cb); + return () => s.stateListeners.delete(cb); + }, + onProgress: (cb) => + { + s.progressListeners.add(cb); + return () => s.progressListeners.delete(cb); + }, + cancel: () => { nativeCleanup(s.peerName); }, + }; +} + +function register(s: NativeSess): void +{ + sessionsById.set(s.sessionId, s); + idByPeer.set(s.peerName, s.sessionId); +} + +// nativeStartOutgoing 经 Go 引擎发送 filePath(绝对路径)到 peerName。total 用于进度事件 +// 的分母(store 的 bytesTransferred 仍是 UI 真值源)。 +export function nativeStartOutgoing( + sessionId: string, peerName: string, filePath: string, total: number, +): P2PSession +{ + ensureSubscribed(); + const s: NativeSess = { + sessionId, peerName, role: "sender", state: "connecting", bytes: 0, total, + stateListeners: new Set(), progressListeners: new Set(), p2pPosted: false, + chain: Promise.resolve(), + }; + register(s); + const ice = JSON.stringify(getICEServers()); + // 用 StartOutgoing 的桥调用作信令链 seed:入站 answer / ICE 须等会话建好再投。 + s.chain = nativeP2PStartOutgoing(sessionId, peerName, filePath, ice).catch((e) => + { + // eslint-disable-next-line no-console + console.error("native start outgoing failed", e); + applyState(s, "failed"); + }); + return makeSession(s); +} + +export function nativeStartIncoming(sessionId: string, senderName: string): P2PSession +{ + ensureSubscribed(); + const s: NativeSess = { + sessionId, peerName: senderName, role: "receiver", state: "connecting", bytes: 0, total: 0, + stateListeners: new Set(), progressListeners: new Set(), p2pPosted: false, + chain: Promise.resolve(), + }; + register(s); + const ice = JSON.stringify(getICEServers()); + // 用 StartIncoming 的桥调用作信令链 seed:入站 offer / ICE 须等会话建好再投。 + s.chain = nativeP2PStartIncoming(sessionId, senderName, ice).catch((e) => + { + // eslint-disable-next-line no-console + console.error("native start incoming failed", e); + applyState(s, "failed"); + }); + return makeSession(s); +} + +export function nativeHandleSignal(fromPeer: string, payload: unknown): void +{ + const payloadJSON = JSON.stringify(payload); + const sid = idByPeer.get(fromPeer); + const s = sid ? sessionsById.get(sid) : undefined; + if (!s) + { + // 无登记会话(极少见):直接尽力转发。 + void nativeP2PHandleSignal(fromPeer, payloadJSON).catch(() => { /* Go 会报无会话 */ }); + return; + } + // 串到会话信令链尾:等会话建好、且按到达顺序逐条投给 Go(offer/answer 先于 ICE 候选)。 + s.chain = s.chain.then(() => nativeP2PHandleSignal(fromPeer, payloadJSON)).catch((e) => + { + // eslint-disable-next-line no-console + console.warn("native handle signal failed", e); + }); +} + +export function nativeCleanup(peerName: string): void +{ + const sid = idByPeer.get(peerName); + if (!sid) { return; } + idByPeer.delete(peerName); + const s = sessionsById.get(sid); + sessionsById.delete(sid); + void nativeP2PCancel(sid); + if (s && s.state !== "completed" && s.state !== "failed") + { + s.state = "closed"; + for (const cb of s.stateListeners) { cb("closed"); } + } +} + +// nativeHasPeer 供 p2p.ts 判断某对端的信令 / 清理该路由到原生后端还是 JS 后端。 +export function nativeHasPeer(peerName: string): boolean +{ + return idByPeer.has(peerName); +} diff --git a/web/src/features/transfer/source.ts b/web/src/features/transfer/source.ts index b8bf091..0471c71 100644 --- a/web/src/features/transfer/source.ts +++ b/web/src/features/transfer/source.ts @@ -8,6 +8,10 @@ export interface FileSource readonly name: string; readonly size: number; readonly type: string; + // path 仅在桌面原生数据面存在:经原生文件选择获得的绝对路径,Go 引擎据此直接读盘做 + // P2P(不走 slice)。浏览器 / iOS 的 File / Range 来源无此字段。slice 仍保留,供中继 + // 回退读取(桌面经 Go 桥读盘片,见 net/desktop.nativeFileSource)。 + readonly path?: string; // 读取 [start, end) 字节区间(end 独占),返回该区间的副本。底层 buffer 固定为 // ArrayBuffer(非 SharedArrayBuffer),以满足 RTCDataChannel.send / fetch body 的类型约束。 slice(start: number, end: number): Promise>; diff --git a/web/src/net/desktop.ts b/web/src/net/desktop.ts index 7e4aa6f..b51e31b 100644 --- a/web/src/net/desktop.ts +++ b/web/src/net/desktop.ts @@ -1,6 +1,7 @@ // Wails 桌面壳的桥接层:仅在桌面 WebView 中可用。浏览器里 isDesktop() 返回 // false,所有调用点都回退到既有的 Web 行为,因此本模块对纯 Web 构建零影响 // ——同一份代码同时部署到浏览器与桌面。 +import type { FileSource } from "../features/transfer/source"; import { useAppStore, type User } from "../store"; // 桌面会话视图:不含 refresh_token——长寿命密钥只在 Go 进程,绝不进 JS(凭据策略 @@ -52,6 +53,25 @@ interface DesktopBridge // 弹原生系统通知(macOS UNUserNotificationCenter / Windows toast)。标题、正文 // 已由前端本地化。 ShowNotification: (title: string, body: string) => Promise; + // 原生 P2P 数据面(pion in Go):发起 / 预备接收 / 喂入站信令 / 取消。进度 / 状态 / + // 出站信令 / 落盘路径经 p2p:* 事件反向回来(见 subscribeNativeP2P)。HTTP 全留 JS。 + // 见 desktop/NATIVE-TRANSFER.md。 + P2PStartOutgoing: (sessionId: string, peerName: string, filePath: string, iceServersJSON: string) => Promise; + P2PStartIncoming: (sessionId: string, peerName: string, iceServersJSON: string) => Promise; + P2PHandleSignal: (fromPeer: string, payloadJSON: string) => Promise; + P2PCancel: (sessionId: string) => Promise; + // 原生文件选择(发送端走原生数据面:Go 按返回的绝对路径直接读盘)。用户取消返回 null。 + PickFileForSend: (title: string) => Promise; + // 中继回退时 WebView 读原生文件片(base64);P2P 主路径不经此桥。 + ReadFileSlice: (path: string, start: number, end: number) => Promise; +} + +// 原生文件选择结果,字段与 Go PickedFile 的 json tag 一致。 +export interface PickedFile +{ + path: string; + name: string; + size: number; } // 桌面剪贴板自动同步开关,默认开;由设置页通过 setClipboardSyncEnabled 调整。 @@ -343,3 +363,121 @@ export async function effectiveDownloadDir(): Promise return ""; } } + +// base64ToBytes 是 bytesToBase64 的逆:把 base64 串解回字节(ArrayBuffer 背衬,满足 +// RTCDataChannel.send / fetch body 的类型约束)。中继回退读原生文件片时用。 +function base64ToBytes(b64: string): Uint8Array +{ + const bin = atob(b64); + const out = new Uint8Array(bin.length); + for (let i = 0; i < bin.length; i += 1) { out[i] = bin.charCodeAt(i); } + return out; +} + +// --- 原生 P2P 数据面桥(仅桌面) --- + +// pickFileForSend 打开原生文件对话框,返回选中文件(绝对路径 + 名 + 大小);取消 / 非桌面 +// 返回 null。供发送端走原生数据面(Go 按路径直接读盘)。 +export async function pickFileForSend(title: string): Promise +{ + const app = bridge(); + if (!app) { return null; } + try { return (await app.PickFileForSend(title)) ?? null; } + catch { return null; } +} + +// nativeFileSource 把原生选中的文件包成 FileSource:path 让 Go 引擎直接读盘做 P2P;slice +// 仅在中继回退时用(经 Go ReadFileSlice 读盘片、base64 过桥——罕见路径,不影响 P2P 主路径)。 +export function nativeFileSource(picked: PickedFile, type: string): FileSource +{ + return { + name: picked.name, + size: picked.size, + type: type || "application/octet-stream", + path: picked.path, + async slice(start, end) + { + const app = bridge(); + if (!app) { throw new Error("desktop bridge unavailable"); } + return base64ToBytes(await app.ReadFileSlice(picked.path, start, end)); + }, + }; +} + +export async function nativeP2PStartOutgoing( + sessionId: string, peerName: string, filePath: string, iceServersJSON: string, +): Promise +{ + const app = bridge(); + if (!app) { throw new Error("desktop bridge unavailable"); } + await app.P2PStartOutgoing(sessionId, peerName, filePath, iceServersJSON); +} + +export async function nativeP2PStartIncoming( + sessionId: string, peerName: string, iceServersJSON: string, +): Promise +{ + const app = bridge(); + if (!app) { throw new Error("desktop bridge unavailable"); } + await app.P2PStartIncoming(sessionId, peerName, iceServersJSON); +} + +export async function nativeP2PHandleSignal(fromPeer: string, payloadJSON: string): Promise +{ + const app = bridge(); + if (!app) { throw new Error("desktop bridge unavailable"); } + await app.P2PHandleSignal(fromPeer, payloadJSON); +} + +export async function nativeP2PCancel(sessionId: string): Promise +{ + const app = bridge(); + if (!app) { return; } + try { await app.P2PCancel(sessionId); } + catch { /* 取消路径吞错 */ } +} + +// 原生 P2P 事件回调集合(引擎经 Wails 事件反向通知)。 +export interface NativeP2PHandlers +{ + onProgress: (sessionId: string, bytes: number) => void; + onState: (sessionId: string, state: string) => void; + onSignal: (sessionId: string, toPeer: string, payloadJSON: string) => void; + onSaved: (sessionId: string, path: string) => void; + onIcePair: (sessionId: string, local: string, remote: string) => void; +} + +// subscribeNativeP2P 订阅 p2p:* 事件并分发;返回取消订阅函数。非桌面为 no-op。 +export function subscribeNativeP2P(h: NativeP2PHandlers): () => void +{ + const rt = runtime(); + if (!rt) { return () => { /* no-op */ }; } + const offs = [ + rt.EventsOn("p2p:progress", (...d) => + { + const p = d[0] as { sessionId: string; bytes: number }; + h.onProgress(p.sessionId, p.bytes); + }), + rt.EventsOn("p2p:state", (...d) => + { + const p = d[0] as { sessionId: string; state: string }; + h.onState(p.sessionId, p.state); + }), + rt.EventsOn("p2p:signal", (...d) => + { + const p = d[0] as { sessionId: string; to: string; payload: string }; + h.onSignal(p.sessionId, p.to, p.payload); + }), + rt.EventsOn("p2p:saved", (...d) => + { + const p = d[0] as { sessionId: string; path: string }; + h.onSaved(p.sessionId, p.path); + }), + rt.EventsOn("p2p:icepair", (...d) => + { + const p = d[0] as { sessionId: string; local: string; remote: string }; + h.onIcePair(p.sessionId, p.local, p.remote); + }), + ]; + return () => { for (const off of offs) { off(); } }; +} diff --git a/web/src/routes/index.tsx b/web/src/routes/index.tsx index 6641bc3..a23bd13 100644 --- a/web/src/routes/index.tsx +++ b/web/src/routes/index.tsx @@ -61,7 +61,9 @@ function HomePage() if (!selectedDevice || !selectedFile) { return; } try { - await startOutgoingTransfer(selectedDevice, fileSource(selectedFile)); + // selectedFile 已是 FileSource(桌面原生选择带 path → 走原生数据面;浏览器 + // File 包装无 path → 走 JS 引擎)。见 Composer 的取文件分支。 + await startOutgoingTransfer(selectedDevice, selectedFile); setSelectedFile(null); } catch (e) diff --git a/web/src/store/types.ts b/web/src/store/types.ts index 6b0400f..0f2ce9d 100644 --- a/web/src/store/types.ts +++ b/web/src/store/types.ts @@ -1,3 +1,4 @@ +import type { FileSource } from "../features/transfer/source"; import type { Locale } from "../i18n"; // Auth shape — see FRONTEND_DESIGN.md §6. @@ -174,9 +175,11 @@ export interface AppState // ---- ui slice ---- selectedDevice: string | null; - selectedFile: File | null; + // 选中待发文件统一为 FileSource:浏览器 / iOS 是 File 包装,桌面原生选择带绝对 path + // (Go 引擎据此直接读盘做原生 P2P)。发送端据有无 path 选原生 / JS 数据面。 + selectedFile: FileSource | null; setSelectedDevice: (n: string | null) => void; - setSelectedFile: (f: File | null) => void; + setSelectedFile: (f: FileSource | null) => void; // ---- theme slice ---- // 'system' 时移除 data-theme 让 @media (prefers-color-scheme) 接管;