桌面原生数据面:pion 数据通道取代 WebView WebRTC,逃离 256KB rwnd 与渲染器节流
把桌面 P2P 数据面从 WebView 内 JS WebRTC 移到 Go 进程的 pion,逃离 WebKit/WebView2 写死的 256KB SCTP 接收窗与隐藏窗口的渲染器节流。后端零改动、与现有客户端互通(线协议逐字节对齐)。真机实测原生↔原生 / 原生↔iOS(经有线)稳定 7-10MB/s,对比旧 JS 路径的 1.5MB/s 还塌缩是质变。设计与阶段见 desktop/NATIVE-TRANSFER.md。
- engine 包(纯 Go,仅 pion+stdlib、零 Wails 依赖,供将来 gomobile 复用于 iOS):pion 收发,逐字节对齐 web 引擎线协议(DataChannel cdrop-file ordered、文本帧 meta/done/ack、二进制分片、信令 {type,sdp?,candidate?});SetSCTPMaxReceiveBufferSize 顶到 4MB;mDNS QueryOnly 解析对端 .local 候选实现同内网 host↔host 直连;接收端直接写盘(去 base64 桥 / OPFS);进程内端到端测试 + 诊断日志(候选/ICE 态/选中对/发送停滞)。
- app.go:Bind P2PStartOutgoing / StartIncoming / HandleSignal / Cancel + 原生取文件 PickFileForSend + 中继回退读盘片 ReadFileSlice;引擎回调经 EventsEmit p2p:* 反向通知(进度/状态/信令/落盘路径/候选对)。
- JS p2pBackend 抽象(p2pNative.ts):桌面按对端把会话委派给 Go 引擎,transfer.ts / hub.ts 零改动;HTTP 全留 JS(出站信令 POST、状态机 /p2p·接收端 /done·/fail);监听 p2p:* 事件驱动 store、phase 与 iceStats。
- 发送取文件改原生路径:selectedFile 模型 File→FileSource(桌面带绝对 path → Go 直接读盘走原生;浏览器拖放无 path → 回退 JS)。Composer dropzone 在桌面打开原生文件对话框。
- 完成语义修复:接收端完成不立即拆连接(续发最终 ack、留 SSE 终态 Cancel 拆、60s 兜底、终态不可被后续 close 翻转)——修「iOS→mac 实际成功却被误标失败」。
- 剪贴板暂存路径修复:富文本经 Universal Clipboard 暂存为 .rtfd,pasteboard 纯文本表示有时取到该暂存路径——uploadClipboard 共享守卫拦截 + mac 原生读改从 RTF / RTFD 派生纯文本。
This commit is contained in:
@@ -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()
|
||||
}
|
||||
@@ -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])
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user