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() }