Files
admin f1a00d128e 原生数据面协议层:桌面 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 互通
2026-06-28 19:32:37 +08:00

177 lines
5.9 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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)}
}
// 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()
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(sessionID, iceServersJSON, e.ev.OnLog)
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(sessionID, iceServersJSON, e.ev.OnLog)
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()
}