原生数据面协议层:桌面 engine 进度节流 / pion 日志 + iOS 数据面路由

- 桌面 engine(pion 数据面,承 ab57afd 之后的精修):
  - 新增 logging.go——把 pion 内部日志路由到宿主 OnLog(真机无 stderr,连接失败无从查);仅 ice/mdns 作用域放 Debug、余 Info+、Trace 丢弃,按 session 聚合
  - wire.go / session.go / engine.go:进度回调 ~10Hz 节流(progressEmitThrottleMs,高吞吐下每片一回调打满宿主主线程;终态由 emitProgressNow 强发最终值),与 web store push 同量级;engine_test.go 跟进
- web 数据面路由:p2p.ts 翻 IOS_NATIVE=true(iOS 走原生 libwebrtc 引擎)+ 新增 p2pIos.ts(iOS p2p 后端)+ p2pNative.ts / net/ios.ts 跟进
- 数据面设计文档 NATIVE-TRANSFER.md:U1(gomobile+pion)真机证伪 → 翻案 U2(libwebrtc)的依据与实测证据(§7/§8)
- 线协议三端一致(cdrop-file ordered / meta+chunk(64KB)+done+ack / 16MB-4MB 水位 / ack 追平完成),桌面 pion ↔ iOS libwebrtc ↔ web JS 互通
This commit is contained in:
2026-06-28 19:32:37 +08:00
parent ab57afdec2
commit f1a00d128e
11 changed files with 587 additions and 22 deletions
+9 -2
View File
@@ -52,6 +52,13 @@ func New(cfg Config, ev Events) *Engine {
return &Engine{cfg: cfg, ev: ev, byPeer: make(map[string]*session)}
}
// NewEngine 是 gomobile-friendly 构造器:以 downloadDir 字符串 + Events 接口构造,避开 New 的
// Config 值参数——gomobile 不绑定按值传递的结构体(故 New 在 iOS 头里被 skip)。桌面侧仍用 New;
// iOS Swift 经 xcframework 调本函数。二者构造同一 *Engine,无行为差异。
func NewEngine(downloadDir string, ev Events) *Engine {
return New(Config{DownloadDir: downloadDir}, ev)
}
// SetDownloadDir 更新落地目录(设置页改下载目录时调用)。
func (e *Engine) SetDownloadDir(dir string) {
e.mu.Lock()
@@ -75,7 +82,7 @@ func (e *Engine) StartOutgoing(sessionID, peerName, filePath, iceServersJSON str
if fi.IsDir() {
return fmt.Errorf("%s is a directory", filePath)
}
pc, err := newPeerConnection(iceServersJSON)
pc, err := newPeerConnection(sessionID, iceServersJSON, e.ev.OnLog)
if err != nil {
return err
}
@@ -109,7 +116,7 @@ func (e *Engine) StartOutgoing(sessionID, peerName, filePath, iceServersJSON str
// StartIncoming 作为接收端预备:建 PeerConnection 等对端 offer。收到 DataChannel 后按
// 线协议接收并直接落盘。须在对端 offer 到达前调用(宿主在 transfer:incoming 时预备)。
func (e *Engine) StartIncoming(sessionID, peerName, iceServersJSON string) error {
pc, err := newPeerConnection(iceServersJSON)
pc, err := newPeerConnection(sessionID, iceServersJSON, e.ev.OnLog)
if err != nil {
return err
}
+1 -1
View File
@@ -61,7 +61,7 @@ func (p *peerEvents) OnSaved(_ string, path string) {
}
}
func (p *peerEvents) OnIcePair(_ string, _ string, _ string) {}
func (p *peerEvents) OnLog(line string) { p.t.Log(line) }
func (p *peerEvents) OnLog(line string) { p.t.Log(line) }
// TestEngineTransfer 端到端:两个进程内 pion peer(A 发 / B 收)经 host 候选直连,传一个
// 3MB 已知内容的文件,断言落盘字节与源一致。覆盖线协议(meta/chunk/done/ack)、水位回压、
+63
View File
@@ -0,0 +1,63 @@
package engine
import (
"fmt"
"github.com/pion/logging"
)
// emitLoggerFactory 把 pion 内部日志路由到宿主 OnLog——真机上无 stderr,否则「为何连不上」无从查。
// 只在连接建立相关作用域(ice / mdns)放开 Debug(候选对检查、ping 失败、sendto 权限错误等关键线索),
// 其余作用域仅 Info 及以上,避免数据面(sctp/dtls)刷屏。Trace 一律丢弃(过量)。
type emitLoggerFactory struct {
emit func(string)
sessionID string
}
func (f emitLoggerFactory) NewLogger(scope string) logging.LeveledLogger {
return emitLogger{scope: scope, emit: f.emit, sessionID: f.sessionID}
}
type emitLogger struct {
scope string
sessionID string
emit func(string)
}
// fwd 统一以 "session <id> ..." 起头,与 session.go 的诊断行同格式,便于宿主按 session 聚合日志。
func (l emitLogger) fwd(level, msg string) {
if l.emit != nil {
l.emit(fmt.Sprintf("session %s pion %s/%s: %s", l.sessionID, l.scope, level, msg))
}
}
// debugOn 仅对连接建立作用域放开 Debug,限制日志量。
func (l emitLogger) debugOn() bool { return l.scope == "ice" || l.scope == "mdns" }
func (l emitLogger) Trace(string) {}
func (l emitLogger) Tracef(string, ...interface{}) {}
func (l emitLogger) Debug(msg string) {
if l.debugOn() {
l.fwd("D", msg)
}
}
func (l emitLogger) Debugf(format string, args ...interface{}) {
if l.debugOn() {
l.fwd("D", fmt.Sprintf(format, args...))
}
}
func (l emitLogger) Info(msg string) { l.fwd("I", msg) }
func (l emitLogger) Infof(format string, args ...interface{}) {
l.fwd("I", fmt.Sprintf(format, args...))
}
func (l emitLogger) Warn(msg string) { l.fwd("W", msg) }
func (l emitLogger) Warnf(format string, args ...interface{}) {
l.fwd("W", fmt.Sprintf(format, args...))
}
func (l emitLogger) Error(msg string) { l.fwd("E", msg) }
func (l emitLogger) Errorf(format string, args ...interface{}) {
l.fwd("E", fmt.Sprintf(format, args...))
}
+46 -15
View File
@@ -35,9 +35,10 @@ type session struct {
fileSize int64
// 计数器(跨 goroutine:发送循环 / pion 回调)。
bytesSent atomic.Int64 // 已塞进本地 SCTP 缓冲的字节
received atomic.Int64 // 接收端已写盘字节
acked atomic.Int64 // 接收端经 ack 回传的已收字节(发送端进度/完成依据)
bytesSent atomic.Int64 // 已塞进本地 SCTP 缓冲的字节
received atomic.Int64 // 接收端已写盘字节
acked atomic.Int64 // 接收端经 ack 回传的已收字节(发送端进度/完成依据)
lastEmitNs atomic.Int64 // emitProgress 节流时间戳(纳秒):限频进度回调,护宿主主线程
// 接收端落盘
out *os.File
@@ -67,8 +68,15 @@ func newSession(e *Engine, id, peer, role string, pc *webrtc.PeerConnection) *se
// 接收窗顶到 4MB(远高于 WebKit 写死的 256KB),是吞吐专项的核心——使本端作接收端时
// `吞吐 ≈ rwnd/RTT` 的天花板抬高一个量级。pion 默认用真实 IP 的 host 候选(不做 mDNS
// 混淆),故同内网天然落到 host↔host 直连,避开 mac WKWebView 的 srflx 退化。
func newPeerConnection(iceServersJSON string) (*webrtc.PeerConnection, error) {
func newPeerConnection(sessionID, iceServersJSON string, log func(string)) (*webrtc.PeerConnection, error) {
se := webrtc.SettingEngine{}
// 把 pion 内部日志接到宿主 OnLog(真机诊断「为何 ICE 连不上」:候选对检查 / sendto 权限错误等),
// 带上 sessionID 使宿主可按 session 聚合 / 整体复制。
se.LoggerFactory = emitLoggerFactory{emit: log, sessionID: sessionID}
// 候选收集用默认(全接口、双栈 IPv4+IPv6、不排除链路本地):桌面 pion 与 iOS WebKit 的原始互通
// 实测可同内网 host↔host 直连(双 WiFi 网常经 IPv6)。早先为治 pion-on-iOS 加的接口 / 链路本地 /
// IPv4-only 过滤纯属 iOS 定向,iOS 已退回 JS、不再用本引擎,那些过滤反令桌面侧 IPv6 host 对被裁掉、
// 退中继(重大回归),故移除。
se.SetSCTPMaxReceiveBufferSize(sctpReceiveBuffer)
// 解析对端的 mDNS.local)候选。iOS / Safari 出于隐私只播 mDNS host 候选;pion 默认不
// 解析,便无法与之成 host↔host、退而走 srflx/relay 慢路径(实测 mac→iOS 落到 relay↔srflx、
@@ -187,6 +195,7 @@ func (s *session) streamFile() {
_ = s.dc.SendText(string(doneB))
s.waitDrained()
s.waitAck(s.fileSize)
s.emitProgressNow()
s.setState("completed")
}
@@ -288,6 +297,7 @@ func (s *session) finalize() {
s.out = nil
s.eng.ev.OnSaved(s.id, final)
}
s.emitProgressNow()
s.setState("completed")
}
@@ -352,21 +362,42 @@ func (s *session) emitSignal(toPeer string, sdp *webrtc.SessionDescription, cand
// --- 进度 / 状态 / 收尾 ---
func (s *session) emitProgress() {
var b int64
// progressBytes 计算当前应上报的已传字节:发送端优先用接收端 ack 追平值(无 ack 时回退本地
// 已交付估计),接收端用已写盘字节。
func (s *session) progressBytes() 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
}
return min(s.fileSize, a)
}
} else {
b = s.received.Load()
b := s.bytesSent.Load() - int64(s.dc.BufferedAmount())
if b < 0 {
return 0
}
return b
}
s.eng.ev.OnProgress(s.id, b)
return s.received.Load()
}
// emitProgress 节流上报:按 progressEmitThrottleMs 限频(~10Hz),避免高吞吐下每片一回调把宿主
// 主线程打满(iOS 上经桥 evaluateJavaScript 数百次/秒 → UI 无响应)。CompareAndSwap 保证多 goroutine
// (发送循环 / 接收 OnMessage / ack)并发下单发。终态的最终值由 emitProgressNow 强发,不被节流吞掉。
func (s *session) emitProgress() {
now := time.Now().UnixNano()
last := s.lastEmitNs.Load()
if now-last < int64(progressEmitThrottleMs)*int64(time.Millisecond) {
return
}
if !s.lastEmitNs.CompareAndSwap(last, now) {
return
}
s.eng.ev.OnProgress(s.id, s.progressBytes())
}
// emitProgressNow 无视节流强发一次:收尾时确保最终字节数到达宿主(否则末次进度可能被节流吞掉,
// 进度文案停在 99%)。
func (s *session) emitProgressNow() {
s.lastEmitNs.Store(time.Now().UnixNano())
s.eng.ev.OnProgress(s.id, s.progressBytes())
}
// waitBufferLow 等本地 SCTP 缓冲落到 threshold 以下:靠 OnBufferedAmountLow 事件唤醒,
+4
View File
@@ -28,6 +28,10 @@ const (
lowWatermark = 4 * 1024 * 1024
// 接收端把已收字节回传发送端的节流间隔(毫秒),与 web 引擎 ACK_INTERVAL_MS 一致。
ackIntervalMs = 200
// 进度回调上报的节流间隔(毫秒):高吞吐下每片一回调会把宿主主线程打满——iOS 上每片
// 都经桥 evaluateJavaScript,数百次/秒会令 UI 无响应(gomobile 集成实测)。限到 ~10Hz
// 终态另由 emitProgressNow 强发最终值。与 web 引擎 store push 的 10Hz 节流同量级。
progressEmitThrottleMs = 100
// 发送端等接收端 ack 追平总量的上限(秒):超时仍照常收尾,避免 ack 丢失致悬挂。
ackCompleteTimeoutSec = 30
// pion 接收窗口:远高于 WebKit 写死的 256KB,使 pion 作接收端不被 rwnd × RTT 掐死